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

1134 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-08-13 02:51 -0700

1# This file is part of lsst-resources. 

2# 

3# Developed for the LSST Data Management System. 

4# This product includes software developed by the LSST Project 

5# (https://www.lsst.org). 

6# See the COPYRIGHT file at the top-level directory of this distribution 

7# for details of code ownership. 

8# 

9# Use of this source code is governed by a 3-clause BSD-style 

10# license that can be found in the LICENSE file. 

11 

12from __future__ import annotations 

13 

14import base64 

15import enum 

16import io 

17import json 

18import logging 

19import os 

20import posixpath 

21import random 

22import re 

23import stat 

24import threading 

25import time 

26import uuid 

27import xml.etree.ElementTree as eTree 

28from datetime import UTC, datetime 

29from http import HTTPStatus 

30from typing import Any, BinaryIO 

31 

32try: 

33 from typing import override # Python 3.12+ 

34except ImportError: 

35 from typing_extensions import override # Python 3.11 

36 

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

38 

39try: 

40 import fsspec 

41 from fsspec.spec import AbstractFileSystem 

42except ImportError: 

43 fsspec = None 

44 AbstractFileSystem = type 

45 

46import yaml 

47from astropy import units as u 

48from urllib3 import PoolManager, make_headers 

49from urllib3.response import HTTPResponse 

50from urllib3.util import Retry, Timeout, Url, parse_url 

51 

52from lsst.utils.logging import getLogger 

53from lsst.utils.timer import time_this 

54 

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

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

57 

58 

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

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

61 

62 A path of the form "///a/b/c///../d/e/" would be normalized as "/a/b/d/e". 

63 The returned path is always absolute, i.e. starts by "/" and never 

64 ends by "/" except when the path is exactly "/" and does not contain 

65 "." nor "..". It does not contain consecutive "/" either. 

66 

67 Parameters 

68 ---------- 

69 path : `str`, optional 

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

71 

72 Returns 

73 ------- 

74 url : `str` 

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

76 """ 

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

78 

79 

80def normalize_url(url: str, preserve_scheme: bool = False, preserve_path: bool = True) -> str: 

81 """Normalize a URL so that scheme be 'http' or 'https' and the URL path 

82 is normalized. 

83 

84 Parameters 

85 ---------- 

86 url : `str` 

87 URL to normalize (e.g., 'davs://example.org:1234///path/to//../dir/'). 

88 preserve_scheme : `bool` 

89 If True the scheme of `url` will be preserved. Otherwise the scheme 

90 of the returned normalized URL will be 'http' or 'https'. 

91 preserve_path : `bool` 

92 If True, the path of `url` will be preserved in the returned 

93 normalized URL, otherwise, the returned URL will have '/' as path. 

94 

95 Returns 

96 ------- 

97 url : `str` 

98 Normalized URL (e.g. 'https://example.org:1234/path/to/dir'). 

99 """ 

100 parsed = parse_url(url) 

101 if parsed.scheme is None: 101 ↛ 102line 101 didn't jump to line 102 because the condition on line 101 was never true

102 scheme = "http" 

103 else: 

104 scheme = parsed.scheme if preserve_scheme else parsed.scheme.replace("dav", "http") 

105 path = normalize_path(parsed.path) if preserve_path else "/" 

106 return Url(scheme=scheme, host=parsed.host, port=parsed.port, path=path).url 

107 

108 

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

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

111 

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

113 leaking authorization tokens. 

114 

115 Parameters 

116 ---------- 

117 url : `str` 

118 URL to redact. 

119 

120 Returns 

121 ------- 

122 redacted_url : `str` 

123 For instance, when called with an URL like: 

124 

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

126 

127 the returned value would be: 

128 

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

130 """ 

131 parsed_url = urlparse(url) 

132 redacted_query: list[tuple[str, str]] = [] 

133 for pair in parse_qsl(parsed_url.query): 

134 redacted_query.append((pair[0], "...." if pair[0] == "authz" else pair[1])) 

135 

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

137 return urlunparse(redacted_url) 

138 

139 

140class DavConfig: 

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

142 particular storage endpoint. 

143 

144 Parameters 

145 ---------- 

146 config : `dict[str, str]` 

147 Dictionary of configurable settings for the webdav endpoint which 

148 base URL is `config["base_url"]`. 

149 

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

151 

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

153 

154 any object of class `DavResourcePath` like 

155 

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

157 

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

159 """ 

160 

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

162 # server. 

163 DEFAULT_TIMEOUT_CONNECT: float = 10.0 

164 

165 # Timeout in seconds to read the response to a request sent to a server. 

166 # This is total time for reading both the headers and the response body. 

167 # It must be large enough to allow for upload and download of files 

168 # of typical size the webdav client supports. 

169 DEFAULT_TIMEOUT_READ: float = 300.0 

170 

171 # Maximum number of network connections to persist against a single 

172 # "host:port" pair. If this endpoint client needs to issue more 

173 # simultaneous requests than this number, additional network connections 

174 # will be created but won't be persisted after use. 

175 DEFAULT_PERSISTENT_CONNECTIONS_PER_HOST: int = 20 

176 

177 # Size of the buffer (in mebibytes, i.e. 1024*1024 bytes) the webdav 

178 # client of this endpoint will use when sending requests and receiving 

179 # responses. 

180 DEFAULT_BUFFER_SIZE: int = 5 

181 

182 # Size of the block (in mebibytes, i.e. 1024*1024 bytes) the webdav 

183 # client of this endpoint will use for making partial reads. Each partial 

184 # read will request at least this number of bytes, unless the total size 

185 # of the file is lower than this value. 

186 DEFAULT_BLOCK_SIZE: int = 1 

187 

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

189 # under certain conditions. 

190 DEFAULT_RETRIES: int = 3 

191 

192 # Minimal and maximal retry backoff (in seconds) for the client to compute 

193 # the wait time before retrying a request. 

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

195 # every time a request is retried. 

196 DEFAULT_RETRY_BACKOFF_MIN: float = 1.0 

197 DEFAULT_RETRY_BACKOFF_MAX: float = 3.0 

198 

199 # Path to a directory or certificate bundle file where the certificates 

200 # of the trusted certificate authorities can be found. 

201 # Those certificates will be used by the client of the webdav endpoint 

202 # to verify the server's host certificate. 

203 # If None, the certificates trusted by the system are used. 

204 DEFAULT_TRUSTED_AUTHORITIES: str | None = None 

205 

206 # User name and password for the client to authenticate to the server. 

207 # If specified, HTTP basic authentication is used on all requests. 

208 DEFAULT_USER_NAME: str | None = None 

209 DEFAULT_USER_PASSWORD: str | None = None 

210 

211 # Path to the client certificate and associated private key the webdav 

212 # client must present to the server for authentication purposes. 

213 # If None, no client certificate is presented. 

214 DEFAULT_USER_CERT: str | None = None 

215 DEFAULT_USER_KEY: str | None = None 

216 

217 # Token the webdav client must sent to the server for authentication 

218 # purposes. The token may be the value of the token itself or the path 

219 # to a file where the token can be found. 

220 DEFAULT_TOKEN: str | None = None 

221 

222 # If this option is set to True, the webdav client attempts to reuse 

223 # the network connection to the server as long as possible. Note that 

224 # the server can unitaleraly decide to close the connection. 

225 # If disabled, the connection is closed after each request. 

226 DEFAULT_REUSE_CONNECTION: bool = True 

227 

228 # Default checksum algorithm to request the server to compute on every 

229 # file upload. Not al servers support this. 

230 # See RFC 3230 for details. 

231 DEFAULT_REQUEST_CHECKSUM: str | None = None 

232 

233 # If this option is set to True, the webdav client can return objects 

234 # compliant to the fsspec specification. 

235 # See: https://filesystem-spec.readthedocs.io 

236 DEFAULT_ENABLE_FSSPEC: bool = True 

237 

238 # If this option is set to True, memory usage is computed and reported 

239 # when executing in debug mode. Computing memory usage is costly, so only 

240 # set this when debugging. 

241 DEFAULT_COLLECT_MEMORY_USAGE: bool = False 

242 

243 # Accepted checksum algorithms. Must be lowercase. 

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

245 

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

247 if config is None: 

248 config = {} 

249 

250 if (base_url := expand_vars(config.get("base_url"))) is None: 

251 self._base_url = "_default_" 

252 else: 

253 self._base_url = normalize_url(base_url, preserve_path=False) 

254 

255 self._timeout_connect: float = float(config.get("timeout_connect", DavConfig.DEFAULT_TIMEOUT_CONNECT)) 

256 self._timeout_read: float = float(config.get("timeout_read", DavConfig.DEFAULT_TIMEOUT_READ)) 

257 self._persistent_connections_per_host: int = int( 

258 config.get( 

259 "persistent_connections_per_host", 

260 DavConfig.DEFAULT_PERSISTENT_CONNECTIONS_PER_HOST, 

261 ) 

262 ) 

263 self._buffer_size: int = 1_048_576 * int(config.get("buffer_size", DavConfig.DEFAULT_BUFFER_SIZE)) 

264 self._block_size: int = 1_048_576 * int(config.get("block_size", DavConfig.DEFAULT_BLOCK_SIZE)) 

265 self._retries: int = int(config.get("retries", DavConfig.DEFAULT_RETRIES)) 

266 self._retry_backoff_min: float = float( 

267 config.get("retry_backoff_min", DavConfig.DEFAULT_RETRY_BACKOFF_MIN) 

268 ) 

269 self._retry_backoff_max: float = float( 

270 config.get("retry_backoff_max", DavConfig.DEFAULT_RETRY_BACKOFF_MAX) 

271 ) 

272 self._trusted_authorities: str | None = expand_vars( 

273 config.get("trusted_authorities", DavConfig.DEFAULT_TRUSTED_AUTHORITIES) 

274 ) 

275 self._user_name: str | None = expand_vars(config.get("user_name", DavConfig.DEFAULT_USER_NAME)) 

276 self._user_password: str | None = expand_vars( 

277 config.get("user_password", DavConfig.DEFAULT_USER_PASSWORD) 

278 ) 

279 self._user_cert: str | None = expand_vars(config.get("user_cert", DavConfig.DEFAULT_USER_CERT)) 

280 self._user_key: str | None = expand_vars(config.get("user_key", DavConfig.DEFAULT_USER_KEY)) 

281 self._token: str | None = expand_vars(config.get("token", DavConfig.DEFAULT_TOKEN)) 

282 self._reuse_connection: bool = config.get("reuse_connection", DavConfig.DEFAULT_REUSE_CONNECTION) 

283 self._enable_fsspec: bool = config.get("enable_fsspec", DavConfig.DEFAULT_ENABLE_FSSPEC) 

284 self._frontend_urls: list[str] = self._init_frontend_urls(config=config) 

285 self._collect_memory_usage: bool = config.get( 

286 "collect_memory_usage", DavConfig.DEFAULT_COLLECT_MEMORY_USAGE 

287 ) 

288 self._request_checksum: str | None = config.get( 

289 "request_checksum", DavConfig.DEFAULT_REQUEST_CHECKSUM 

290 ) 

291 if self._request_checksum is not None: 

292 self._request_checksum = self._request_checksum.lower() 

293 if self._request_checksum not in DavConfig.ACCEPTED_CHECKSUMS: 293 ↛ 294line 293 didn't jump to line 294 because the condition on line 293 was never true

294 raise ValueError( 

295 f"""Value for checksum algorithm {self._request_checksum} for storage endpoint """ 

296 f"""{self._base_url} is not among the accepted values: {DavConfig.ACCEPTED_CHECKSUMS}""" 

297 ) 

298 

299 def _init_frontend_urls(self, config: dict | None = None) -> list[str]: 

300 if config is None: 300 ↛ 301line 300 didn't jump to line 301 because the condition on line 300 was never true

301 return [] 

302 

303 # Initialize the URLs of the frontend servers, if present in 

304 # the configuration. 

305 frontend_urls: list[str] = [] 

306 for url in config.get("frontend_base_urls", []): 

307 # Expand environment variables in this URL 

308 if (expanded_url := expand_vars(url)) is not None: 308 ↛ 306line 308 didn't jump to line 306 because the condition on line 308 was always true

309 frontend_urls.append(normalize_url(expanded_url, preserve_path=False)) 

310 

311 # Eliminate duplicate URLs. 

312 frontend_urls = list(set(frontend_urls)) 

313 

314 # Check that the scheme of this client's base URL is identical to 

315 # the scheme of the frontend server URLs. 

316 base_url_scheme = parse_url(self._base_url).scheme 

317 for url in frontend_urls: 

318 if base_url_scheme != parse_url(url).scheme: 

319 raise ValueError( 

320 f"""inconsistent scheme in frontend URL {url} for endpoint """ 

321 f"""with base URL {self._base_url}""" 

322 ) 

323 

324 return frontend_urls 

325 

326 @property 

327 def base_url(self) -> str: 

328 return self._base_url 

329 

330 @property 

331 def timeout_connect(self) -> float: 

332 return self._timeout_connect 

333 

334 @property 

335 def timeout_read(self) -> float: 

336 return self._timeout_read 

337 

338 @property 

339 def persistent_connections_per_host(self) -> int: 

340 return self._persistent_connections_per_host 

341 

342 @property 

343 def buffer_size(self) -> int: 

344 return self._buffer_size 

345 

346 @property 

347 def block_size(self) -> int: 

348 return self._block_size 

349 

350 @property 

351 def retries(self) -> int: 

352 return self._retries 

353 

354 @property 

355 def retry_backoff_min(self) -> float: 

356 return self._retry_backoff_min 

357 

358 @property 

359 def retry_backoff_max(self) -> float: 

360 return self._retry_backoff_max 

361 

362 @property 

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

364 return self._trusted_authorities 

365 

366 @property 

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

368 return self._token 

369 

370 @property 

371 def reuse_connection(self) -> bool: 

372 return self._reuse_connection 

373 

374 @property 

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

376 return self._request_checksum 

377 

378 @property 

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

380 return self._user_cert 

381 

382 @property 

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

384 # If no user certificate was specified in the configuration, 

385 # ignore the private key, even if it was provided. 

386 if self._user_cert is None: 

387 return None 

388 

389 # If we have a user certificate but not a private key, assume the 

390 # private key is included in the same file as the user certificate. 

391 # That is typically the case when using a X.509 grid proxy as 

392 # client certificate. 

393 return self._user_cert if self._user_key is None else self._user_key 

394 

395 @property 

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

397 return self._user_name 

398 

399 @property 

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

401 return self._user_password 

402 

403 @property 

404 def enable_fsspec(self) -> bool: 

405 return self._enable_fsspec 

406 

407 @property 

408 def collect_memory_usage(self) -> bool: 

409 return self._collect_memory_usage 

410 

411 @property 

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

413 return self._frontend_urls 

414 

415 

416class DavConfigPool: 

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

418 

419 Parameters 

420 ---------- 

421 filename : `list` [ `str` ] 

422 List of environment variables or file names to load the configuration 

423 from. The first file found in the list will be read and the 

424 configuration settings for all webDAV endpoints will be extracted 

425 from it. Other files will be ignored. 

426 

427 Each component of `filenames` can be an environment variable or 

428 the path of a file which itself can include an environment variable, 

429 e.g. '$HOME/path/to/config.yaml'. 

430 

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

432 

433 - base_url: "davs://webdav1.example.org:1234/" 

434 persistent_connections_per_host: 10 

435 timeout_connect: 20.0 

436 timeout_read: 120.0 

437 retries: 3 

438 retry_backoff_min: 1.0 

439 retry_backoff_max: 3.0 

440 user_cert: "${X509_USER_PROXY}" 

441 user_key: "${X509_USER_PROXY}" 

442 token: "/path/to/bearer/token/file" 

443 trusted_authorities: "/etc/grid-security/certificates" 

444 buffer_size: 5 

445 enable_fsspec: false 

446 request_checksum: "md5" 

447 collect_memory_usage: false 

448 

449 - base_url: "davs://webdav2.example.org:1234/" 

450 user_name: "user" 

451 user_password: "password" 

452 persistent_connections_per_host: 5 

453 reuse_connection: false 

454 ... 

455 

456 All settings are optional. If no settings are found in the 

457 configuration file for a particular webDAV endpoint, sensible 

458 defaults will be used. 

459 

460 There is only a single instance of this class. This thead-safe 

461 singleton is intended to be initialized when the module is imported 

462 the first time. 

463 """ 

464 

465 _instance = None 

466 _lock = threading.Lock() 

467 

468 def __new__(cls, filename: str | None = None) -> DavConfigPool: 

469 if cls._instance is None: 

470 with cls._lock: 

471 if cls._instance is None: 471 ↛ 474line 471 didn't jump to line 474

472 cls._instance = super().__new__(cls) 

473 

474 return cls._instance 

475 

476 def __init__(self, filename: str | None = None) -> None: 

477 # Create a default configuration. This configuration is 

478 # used when a URL doest not match any of the endpoints in the 

479 # configuration. 

480 self._default_config: DavConfig = DavConfig() 

481 

482 # The key of this dictionary is the URL of the webDAV endpoint, 

483 # e.g. "davs://host.example.org:1234/" 

484 self._configs: dict[str, DavConfig] = {} 

485 

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

487 # if any. 

488 if filename is None: 

489 return 

490 

491 # filename can be the name of an environment variable or a path. 

492 # A path can include environment variables 

493 # (e.g. "$HOME/path/to/config.yaml") or "~" 

494 # (e.g. "~/path/to/config.yaml") 

495 if (filename := os.getenv(filename)) is not None: 

496 # Expand environment variables and '~' in the file name, if any. 

497 filename = os.path.expandvars(filename) 

498 filename = os.path.expanduser(filename) 

499 with open(filename) as file: 

500 for config_item in yaml.safe_load(file): 

501 config = DavConfig(config_item) 

502 if config.base_url not in self._configs: 

503 self._configs[config.base_url] = config 

504 else: 

505 # We already have a configuration for the same 

506 # endpoint. That is likely a human error in 

507 # the configuration file. 

508 raise ValueError( 

509 f"""configuration file {filename} contains two configurations for """ 

510 f"""endpoint {config.base_url}""" 

511 ) 

512 

513 def get_config_for_url(self, url: str) -> DavConfig: 

514 """Return the configuration to use a webDAV client when interacting 

515 with the server which hosts the resource at `url`. 

516 

517 Parameters 

518 ---------- 

519 url : `str` 

520 URL for which to obtain a configuration. 

521 """ 

522 # Select the configuration for the endpoint of the provided URL. 

523 normalized_url: str = normalize_url(url, preserve_path=False) 

524 if (config := self._configs.get(normalized_url)) is not None: 

525 return config 

526 

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

528 return self._default_config 

529 

530 def _destroy(self) -> None: 

531 """Destroy this class singleton instance. 

532 

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

534 """ 

535 with DavConfigPool._lock: 

536 DavConfigPool._instance = None 

537 

538 

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

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

541 

542 Parameters 

543 ---------- 

544 config : `DavConfig` 

545 Configurable settings for a webDAV storage endpoint. 

546 

547 Returns 

548 ------- 

549 retry : `urllib3.util.Retry` 

550 Retry object to he used when creating a ``urllib3.PoolManager``. 

551 """ 

552 backoff_min: float = config.retry_backoff_min 

553 backoff_max: float = config.retry_backoff_max 

554 retry = Retry( 

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

556 # counts. 

557 total=3 * config.retries, 

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

559 connect=config.retries, 

560 # How many times to retry on read errors. 

561 read=config.retries, 

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

563 status=config.retries, 

564 # How many times to retry on other errors. 

565 other=config.retries, 

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

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

568 # started at the same time (even on different hosts) to overwhelm the 

569 # server by sending requests at the same time. 

570 backoff_factor=backoff_min + (backoff_max - backoff_min) * random.random(), 

571 # How many redirects to perform. Set to a finite value to avoid 

572 # infinite redirect loops. 

573 redirect=3, 

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

575 # By default, we automatically retry idempotent requests. Specific 

576 # retry configuration may be set for some requests such as 

577 # non-idempotent `PUT` requests. 

578 allowed_methods=frozenset( 

579 [ 

580 "COPY", 

581 "DELETE", 

582 "GET", 

583 "HEAD", 

584 "MKCOL", 

585 "OPTIONS", 

586 "PROPFIND", 

587 "PUT", 

588 ] 

589 ), 

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

591 status_forcelist=frozenset( 

592 [ 

593 HTTPStatus.TOO_MANY_REQUESTS, # 429 

594 HTTPStatus.INTERNAL_SERVER_ERROR, # 500 

595 HTTPStatus.BAD_GATEWAY, # 502 

596 HTTPStatus.SERVICE_UNAVAILABLE, # 503 

597 HTTPStatus.GATEWAY_TIMEOUT, # 504 

598 ] 

599 ), 

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

601 # above. 

602 respect_retry_after_header=True, 

603 ) 

604 return retry 

605 

606 

607class DavClientPool: 

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

609 to talk to a single storage endpoint. 

610 

611 Parameters 

612 ---------- 

613 config_pool : `DavConfigPool` 

614 Pool of all known webDAV client configurations. 

615 

616 Notes 

617 ----- 

618 There is a single instance of this class. This thead-safe singleton is 

619 intended to be initialized when the module is imported the first time. 

620 """ 

621 

622 _instance = None 

623 _lock = threading.Lock() 

624 

625 def __new__(cls, config_pool: DavConfigPool) -> DavClientPool: 

626 if cls._instance is None: 626 ↛ 631line 626 didn't jump to line 631 because the condition on line 626 was always true

627 with cls._lock: 

628 if cls._instance is None: 628 ↛ 631line 628 didn't jump to line 631

629 cls._instance = super().__new__(cls) 

630 

631 return cls._instance 

632 

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

634 self._config_pool: DavConfigPool = config_pool 

635 

636 # The key of this dictionnary is a path-stripped URL of the form 

637 # "davs://host.example.org:1234/". The value is a reusable 

638 # DavClient to interact with that endpoint. 

639 self._clients: dict[str, DavClient] = {} 

640 

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

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

643 is hosted. 

644 

645 Parameters 

646 ---------- 

647 url : `str` 

648 URL for which to obtain a client. 

649 

650 Notes 

651 ----- 

652 The returned client is thread-safe. If a client for that endpoint 

653 already exists it is reused, otherwise a new client is created 

654 with the appropriate configuration for interacting with the storage 

655 endpoint. 

656 """ 

657 # If we already have a client for this endpoint reuse it. 

658 url = normalize_url(url, preserve_path=False) 

659 if (client := self._clients.get(url)) is not None: 

660 return client 

661 

662 # No client for this endpoint was found. Create a new one and save it 

663 # for serving subsequent requests. 

664 with DavClientPool._lock: 

665 # If another client was created in the meantime by another thread 

666 # reuse it. 

667 if (client := self._clients.get(url)) is not None: 

668 return client 

669 

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

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

672 

673 return self._clients[url] 

674 

675 def _make_client(self, url: str, config: DavConfig) -> DavClient: 

676 """Make a webDAV client for interacting with the server at `url`.""" 

677 # Check the server implements webDAV protocol and retrieve its 

678 # identity so that we can build a client for that specific 

679 # server implementation. 

680 client = DavClient(url, config) 

681 server_details = client.get_server_details(url) 

682 server_id = server_details.get("Server", None) 

683 accepts_ranges: bool | str | None = server_details.get("Accept-Ranges", None) 

684 if accepts_ranges is not None: 

685 accepts_ranges = accepts_ranges == "bytes" 

686 

687 if server_id is None: 

688 # Create a generic webDAV client 

689 return DavClient(url, config, accepts_ranges) 

690 server_id = server_id.lower() 

691 if server_id.startswith("dcache"): 

692 # Create a client for a dCache webDAV server 

693 return DavClientDCache(url, config, accepts_ranges) 

694 elif server_id.startswith("xrootd"): 

695 # Create a client for a XrootD webDAV server 

696 return DavClientXrootD(url, config, accepts_ranges) 

697 else: 

698 # Return a generic webDAV client 

699 return DavClient(url, config, accepts_ranges) 

700 

701 def _destroy(self) -> None: 

702 """Destroy this class singleton instance. 

703 

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

705 """ 

706 with DavClientPool._lock: 

707 DavClientPool._instance = None 

708 

709 

710class DavFileSizeCache: 

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

712 

713 Parameters 

714 ---------- 

715 default_timeout : `float`, optional 

716 Default validity period, in seconds, of the entries in this cache. 

717 The validity period for a specific entry can be specified when the 

718 entry is added to the cache (see `update_size` method). 

719 

720 Notes 

721 ----- 

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

723 objects. This singleton is thread safe. 

724 

725 Caching file sizes helps preventing sending requests to the server for 

726 retrieving the size of recently uploaded files. This is in particular 

727 intended to efficiently serve `Butler` requests for the size of a file it 

728 just wrote to the datastore. 

729 """ 

730 

731 _instance = None 

732 _lock = threading.Lock() 

733 

734 def __new__(cls) -> DavFileSizeCache: 

735 if cls._instance is None: 

736 with cls._lock: 

737 if cls._instance is None: 

738 cls._instance = super().__new__(cls) 

739 

740 return cls._instance 

741 

742 def __init__(self, default_timeout: float = 60.0) -> None: 

743 # The key of the cache dictionnary is a URL of the form 

744 # 

745 # "https://host.example.org:1234/path/to/file". 

746 # 

747 # The value is a triplet (file_size, last_updated, timeout) where: 

748 # - 'file_size' is the size of the file in bytes, 

749 # - 'last_updated' is the time when this entry was added to the cache 

750 # or last updated, in seconds since epoch, 

751 # - 'timeout' is the validity period of this cache entry, in seconds, 

752 # understood from the moment the cache entry was created. 

753 with DavFileSizeCache._lock: 

754 if not hasattr(self, "_cache"): 

755 self._default_timeout: float = default_timeout 

756 self._cache: dict[str, tuple[int, float, float]] = {} 

757 

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

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

760 

761 Parameters 

762 ---------- 

763 url : `str` 

764 URL of the file to invalidate which cache entry must be 

765 invalidated. 

766 """ 

767 with DavFileSizeCache._lock: 

768 self._cache.pop(url, None) 

769 

770 def update_size(self, url: str, size: int | None, timeout: float | None = None) -> None: 

771 """Update the cache with an entry for `url` which has a size of `size` 

772 bytes. This entry is considered valid for a period of `timeout` 

773 seconds from now. 

774 

775 Parameters 

776 ---------- 

777 url : `str` 

778 URL of the file the size to be cached. 

779 size : `size` or `None`, optional 

780 Size in bytes of the file at `url`. If this value is `None`, the 

781 cache is not modified. 

782 timeout : `float` or `None`, optional 

783 The validity period, in seconds, this size is to be considered 

784 valid. If not specified, the default value specified when this 

785 object was created will be used for this cache entry. 

786 """ 

787 if size is None: 

788 return 

789 

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

791 with DavFileSizeCache._lock: 

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

793 

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

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

796 

797 Parameters 

798 ---------- 

799 url : `str` 

800 URL of the file to retrieve the size for. 

801 

802 Returns 

803 ------- 

804 `size`: `int` or `None` 

805 The cached value of the size of file at `url` if any value was 

806 found in the cache, `None` otherwise. 

807 `None` is also returned if there is a cached value but its 

808 validity period has expired. In this case, the entry associated to 

809 `url` is removed from the cache. 

810 """ 

811 with DavFileSizeCache._lock: 

812 if (entry := self._cache.get(url, None)) is None: 

813 # There is no entry in the cache for this URL 

814 return None 

815 

816 # There is an entry in the cache for this URL. Check that 

817 # its validity period has not yet expired. 

818 size, last_updated, timeout = entry 

819 if time.time() <= last_updated + timeout: 

820 # This entry is stil valid 

821 return size 

822 else: 

823 # This entry is no longer valid. Remove it from the cache. 

824 self._cache.pop(url) 

825 return None 

826 

827 

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

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

830 

831 Parameters 

832 ---------- 

833 method : `str` 

834 The method name triggering the error. 

835 url : `str` 

836 The URL that cause the error. 

837 resp : `resp` 

838 The error response. 

839 """ 

840 message = f"Unexpected response to HTTP request {method} {redact_url(url)}: {resp.status} {resp.reason}" 

841 body = resp.data.decode() 

842 if len(body) > 0: 

843 message += f" [response body: {body}]" 

844 

845 return ValueError(message) 

846 

847 

848class DavClient: 

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

850 

851 Instances of this class are thread-safe. 

852 

853 Parameters 

854 ---------- 

855 url : `str` 

856 Root URL of the storage endpoint (e.g. 

857 "https://host.example.org:1234/"). 

858 config : `DavConfig` 

859 Configuration to initialize this client. 

860 accepts_ranges : `bool` | `None` 

861 Indicate whether the remote server accepts the ``Range`` header in GET 

862 requests. 

863 """ 

864 

865 def __init__(self, url: str, config: DavConfig, accepts_ranges: bool | None = None) -> None: 

866 # Lock to protect this client fields from concurrent modification. 

867 self._lock = threading.Lock() 

868 

869 # Base URL of the server this client will interact with. 

870 # It is of the form: "davs://host.example.org:1234/" 

871 self._base_url: str = url 

872 

873 # Configuration settings for the storage endpoint this client 

874 # will interact with. 

875 self._config: DavConfig = config 

876 

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

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

879 

880 # Make the pool manager for this client to use for sending 

881 # requests to the server. 

882 self._pool_manager: PoolManager = self._make_pool_manager(config=self._config) 

883 

884 # Parser of PROPFIND responses. 

885 self._propfind_parser: DavPropfindParser = DavPropfindParser() 

886 

887 # Does the remote server accept a "Range" header in GET requests? 

888 # This field is lazy initialized. 

889 self._accepts_ranges: bool | None = accepts_ranges 

890 

891 # Can this client use a COPY request to duplicate files within a 

892 # single webDAV server? 

893 # Subclasses can overwrite this setting according to the server 

894 # capabilities and compliance to webDAV RFC. 

895 self._can_duplicate: bool = True 

896 

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

898 # to the server. 

899 self._file_size_cache = DavFileSizeCache() 

900 

901 def _make_authorizer(self, config: DavConfig) -> Authorizer | None: 

902 # If a token was specified in the configuration settings for this 

903 # endpoint, prefer it as the authentication method, even if other 

904 # authentication settings were also specified. 

905 if config.token is not None: 

906 return TokenAuthorizer(token=config.token) 

907 elif config.user_name is not None and config.user_password is not None: 

908 return BasicAuthorizer(user_name=config.user_name, user_password=config.user_password) 

909 

910 return None 

911 

912 def _make_pool_manager(self, config: DavConfig) -> PoolManager: 

913 # Prepare the trusted authorities certificates 

914 ca_certs, ca_cert_dir = None, None 

915 if config.trusted_authorities is not None: 

916 if os.path.isdir(config.trusted_authorities): 

917 ca_cert_dir = config.trusted_authorities 

918 elif os.path.isfile(config.trusted_authorities): 

919 ca_certs = config.trusted_authorities 

920 else: 

921 raise FileNotFoundError( 

922 f"Trusted authorities file or directory {config.trusted_authorities} does not exist" 

923 ) 

924 

925 # If a token was specified for this endpoint don't use the 

926 # <user certificate, private key> pair, even if they were also 

927 # specified. 

928 user_cert, user_key = None, None 

929 if config.token is None: 

930 user_cert = config.user_cert 

931 user_key = config.user_key 

932 

933 # Pool manager for sending requests. Connections in this pool manager 

934 # are generally left open by the client but the front-end server may 

935 # choose to close them in some specific situations. For instance, 

936 # whe serving a PUT request, the front server may redirect to a 

937 # backend server and close the network connection making it 

938 # unsuable for subsequent requests. 

939 # 

940 # In addition, the client may also choose to explicitly close the 

941 # network connection after receiving a response. 

942 return PoolManager( 

943 # Number of connection pools to cache before discarding the least 

944 # recently used pool. Each connection pool manages network 

945 # connections to a single host, so this is basically the number 

946 # of "host:port" we persist network connections to. 

947 num_pools=200, 

948 # Number of connections to the same "host:port" to persist for 

949 # later reuse. More than 1 is useful in multithreaded situations. 

950 # If more than this number of network connections are needed at 

951 # a particular moment, they will be created and discarded after 

952 # use. 

953 maxsize=config.persistent_connections_per_host, 

954 # Retry configuration to use by default with requests sent to 

955 # host in the front end. 

956 retries=make_retry(config), 

957 # Socket timeout in seconds for each individual connection. 

958 timeout=Timeout( 

959 connect=config.timeout_connect, 

960 read=config.timeout_read, 

961 ), 

962 # Size in bytes of the buffer for reading/writing data from/to 

963 # the underlying socket. 

964 blocksize=config.buffer_size, 

965 # Client certificate and private key for esablishing TLS 

966 # connections. If None, no client certificate is sent to the 

967 # server. Only relevant for endpoints using secure HTTP protocol. 

968 cert_file=user_cert, 

969 key_file=user_key, 

970 # We require verification of the server certificate. 

971 cert_reqs="CERT_REQUIRED", 

972 # Directory where the certificates of the trusted certificate 

973 # authorities can be found. The contents of that directory 

974 # must be as expected by OpenSSL. 

975 ca_cert_dir=ca_cert_dir, 

976 # Path to a file of concatenated CA certificates in PEM format. 

977 ca_certs=ca_certs, 

978 ) 

979 

980 def get_server_details(self, url: str) -> dict[str, str]: 

981 """Retrieve the details of the server and check it advertises 

982 compliance to class 1 of webDAV protocol. 

983 

984 Parameters 

985 ---------- 

986 url : `str` 

987 URL to check. 

988 

989 Returns 

990 ------- 

991 details: `dic[str, str]` 

992 The keys of the returned dictionary can be "Server" and 

993 "Accept-Ranges". Any of those keys may not exist in the returned 

994 dictionary if the server did not include it in its response. 

995 

996 The values are the values of the corresponding 

997 headers found in the response to the OPTIONS request. 

998 Examples of values for the "Server" header are 'dCache/9.2.4' or 

999 'XrootD/v5.7.1'. 

1000 """ 

1001 # Check that the value "1" is part of the value of the "DAV" header in 

1002 # the response to an 'OPTIONS' request. 

1003 # 

1004 # We don't rely on webDAV locks, so a server complying to class 1 is 

1005 # enough for our purposes. All webDAV servers must advertise at least 

1006 # compliance class "1". 

1007 # 

1008 # Compliance classes are documented in 

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

1010 # 

1011 # Examples of values for header DAV are: 

1012 # DAV: 1, 2 

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

1014 resp = self.options(url) 

1015 if "DAV" not in resp.headers: 

1016 raise ValueError(f"Server of {resp.geturl()} does not implement webDAV protocol") 

1017 

1018 if "1" not in resp.headers.get("DAV").replace(" ", "").split(","): 

1019 raise ValueError( 

1020 f"Server of {resp.geturl()} does not advertise required compliance to webDAV protocol class 1" 

1021 ) 

1022 

1023 # The value of 'Server' header is expected to be of the form 

1024 # 'dCache/9.2.4' or 'XrootD/v5.7.1'. Not all servers include such a 

1025 # header in their response to an OPTIONS request. 

1026 details: dict[str, str] = {} 

1027 for header in ("Server", "Accept-Ranges"): 

1028 value = resp.headers.get(header, None) 

1029 if value is not None: 

1030 details[header] = value 

1031 

1032 return details 

1033 

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

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

1036 

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

1038 """ 

1039 if resp.retries is None: 

1040 return default_url 

1041 

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

1043 return default_url 

1044 

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

1046 

1047 def _rewrite_url_for_frontend(self, url: str) -> str: 

1048 """Return a URL to reach one of the frontend servers that serves 

1049 requests sent against `url`. 

1050 

1051 Parameters 

1052 ---------- 

1053 url : `str` 

1054 Target URL. 

1055 

1056 Returns 

1057 ------- 

1058 url: `str` 

1059 URL to reach one of this client's frontend servers. If `url` does 

1060 not target this client's frontend servers, the returned value 

1061 is `url` unmodified. 

1062 """ 

1063 # Do nothing if this URL does not match this client's base URL. This 

1064 # happens, for instance, when `url` is a redirection to a backend 

1065 # server, so we don't want to rewrite it. 

1066 # 

1067 # Also, don't rewrite the URL if we don't have frontends configured 

1068 # for this client. 

1069 if not self._config.frontend_urls or not url.startswith(self._base_url): 

1070 return url 

1071 

1072 # Randomly select one of the configured frontends and return a modified 

1073 # URL which uses the selected frontend instead of the original one. 

1074 return random.choice(self._config.frontend_urls) + url.removeprefix(self._base_url) 

1075 

1076 def _request( 

1077 self, 

1078 method: str, 

1079 url: str, 

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

1081 body: BinaryIO | bytes | str | None = None, 

1082 preload_content: bool = True, 

1083 redirect: bool = True, 

1084 pool_manager: PoolManager | None = None, 

1085 **kwargs: dict[Any, Any], 

1086 ) -> HTTPResponse: 

1087 """Send a generic HTTP request and return the response. 

1088 

1089 Parameters 

1090 ---------- 

1091 method : `str` 

1092 Request method, e.g. 'GET', 'PUT', 'PROPFIND'. 

1093 url : `str` 

1094 Target URL. 

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

1096 Headers to sent with the request. 

1097 body : `bytes` or `str` or `None`, optional 

1098 Request body. 

1099 preload_content : `bool`, optional 

1100 If True, the response body is downloaded and can be retrieved 

1101 via the returned response `.data` property. If False, the 

1102 caller needs to call `.read()` on the returned response object to 

1103 download the body, either entirely in one call or by chunks. 

1104 redirect : `bool`, optional 

1105 If True, automatically handle redirects. If False, the returned 

1106 response may contain a redirection to another location. 

1107 pool_manager : `PoolManager`, optional 

1108 Pool manager to use for sending this request. If not provided, 

1109 this client's pool manager is used. 

1110 kwargs : `dict[Any, Any]`, optional 

1111 Keyword arguments to pass unmodified to 

1112 `urllib3.PoolManager.request()`. 

1113 

1114 Returns 

1115 ------- 

1116 resp: `HTTPResponse` 

1117 Response to the request as received from the server. 

1118 """ 

1119 # Retrieve the URL we must use to send this request to one of this 

1120 # client's configured frontend servers. 

1121 url = self._rewrite_url_for_frontend(url) 

1122 

1123 # If this client is configured not to reuse the network connection 

1124 # with the server, add a "Connection: close" header to this request. 

1125 # 

1126 # However, if the caller has explicitly specified a "Connection" 

1127 # header, whatever its value, don't modify it. 

1128 headers = {} if headers is None else dict(headers) 

1129 if "Connection" not in headers and not self._config.reuse_connection: 

1130 headers.update({"Connection": "close"}) 

1131 

1132 # If an authorizer (basic or token) is configured for this client, 

1133 # allow it to set the "Authorization" header to this outgoing request. 

1134 if self._authorizer is not None: 

1135 self._authorizer.set_authorization(headers) 

1136 

1137 if log.isEnabledFor(logging.DEBUG): 

1138 annotation = "" 

1139 if method == "GET" and "Range" in headers: 

1140 byte_range = headers.get("Range", "").removeprefix("bytes=") 

1141 annotation = f" (byte range: {byte_range})" 

1142 

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

1144 

1145 if pool_manager is None: 

1146 pool_manager = self._pool_manager 

1147 

1148 with time_this( 

1149 log, 

1150 msg="%s %s", 

1151 args=(method, url), 

1152 mem_usage=self._config.collect_memory_usage, 

1153 mem_unit=u.mebibyte, 

1154 ): 

1155 return pool_manager.request( 

1156 method, 

1157 url, 

1158 body=body, 

1159 headers=headers, 

1160 preload_content=preload_content, 

1161 redirect=redirect, 

1162 **kwargs, 

1163 ) 

1164 

1165 def _options( 

1166 self, 

1167 url: str, 

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

1169 pool_manager: PoolManager | None = None, 

1170 ) -> HTTPResponse: 

1171 """Send a HTTP OPTIONS request and return the response unmodified. 

1172 

1173 Parameters 

1174 ---------- 

1175 url : `str` 

1176 Target URL. 

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

1178 Headers to sent with the request. 

1179 pool_manager : `PoolManager`, optional 

1180 Pool manager to use to send this request. 

1181 

1182 Returns 

1183 ------- 

1184 resp: `HTTPResponse` 

1185 Response to the request as received from the server. 

1186 

1187 Notes 

1188 ----- 

1189 This method is intended for subclasses to override when needed. 

1190 """ 

1191 return self._request("OPTIONS", url=url, headers=headers, pool_manager=pool_manager) 

1192 

1193 def _copy( 

1194 self, 

1195 url: str, 

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

1197 preload_content: bool = True, 

1198 pool_manager: PoolManager | None = None, 

1199 ) -> HTTPResponse: 

1200 """Send a webDAV COPY request and return the response unmodified. 

1201 

1202 Parameters 

1203 ---------- 

1204 url : `str` 

1205 Target URL. 

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

1207 Headers to sent with the request. 

1208 pool_manager : `PoolManager`, optional 

1209 Pool manager to use to send this request. 

1210 

1211 Notes 

1212 ----- 

1213 This method is intended for subclasses to override when needed. 

1214 """ 

1215 return self._request( 

1216 "COPY", url=url, headers=headers, preload_content=preload_content, pool_manager=pool_manager 

1217 ) 

1218 

1219 def _delete( 

1220 self, 

1221 url: str, 

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

1223 pool_manager: PoolManager | None = None, 

1224 ) -> HTTPResponse: 

1225 """Send a HTTP DELETE request and return the response unmodified. 

1226 

1227 Parameters 

1228 ---------- 

1229 url : `str` 

1230 Target URL. 

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

1232 Headers to sent with the request. 

1233 pool_manager : `PoolManager`, optional 

1234 Pool manager to use to send this request. 

1235 

1236 Notes 

1237 ----- 

1238 This method is intended for subclasses to override when needed. 

1239 """ 

1240 return self._request("DELETE", url=url, headers=headers, pool_manager=pool_manager) 

1241 

1242 def _get( 

1243 self, 

1244 url: str, 

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

1246 preload_content: bool = True, 

1247 redirect: bool = True, 

1248 pool_manager: PoolManager | None = None, 

1249 ) -> HTTPResponse: 

1250 """Send a HTTP GET request and return the response unmodified. 

1251 

1252 Parameters 

1253 ---------- 

1254 url : `str` 

1255 Target URL. 

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

1257 Headers to sent with the request. 

1258 preload_content : `bool`, optional 

1259 If True, the response body is downloaded and can be retrieved 

1260 via the returned response `.data` property. If False, the 

1261 caller needs to call the `.read()` on the returned response 

1262 object to download the body. 

1263 redirect : `bool`, optional 

1264 If True, follow redirections. 

1265 pool_manager : `PoolManager`, optional 

1266 Pool manager to send the request through. 

1267 

1268 Returns 

1269 ------- 

1270 resp: `HTTPResponse` 

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

1272 

1273 Notes 

1274 ----- 

1275 This method is intended for subclasses to override when needed. 

1276 """ 

1277 return self._request( 

1278 "GET", 

1279 url=url, 

1280 headers=headers, 

1281 preload_content=preload_content, 

1282 redirect=redirect, 

1283 pool_manager=pool_manager, 

1284 ) 

1285 

1286 def _head( 

1287 self, 

1288 url: str, 

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

1290 pool_manager: PoolManager | None = None, 

1291 ) -> HTTPResponse: 

1292 """Send a HTTP HEAD request and return the response. 

1293 

1294 Parameters 

1295 ---------- 

1296 url : `str` 

1297 Target URL. 

1298 headers : `bool` 

1299 If the target URL is not found, raise an exception. Otherwise 

1300 just return the response. 

1301 pool_manager : `PoolManager`, optional 

1302 Pool manager to use to send this request. 

1303 

1304 Notes 

1305 ----- 

1306 This method is intended for subclasses to override when needed. 

1307 """ 

1308 return self._request("HEAD", url=url, headers=headers, pool_manager=pool_manager) 

1309 

1310 def _mkcol( 

1311 self, 

1312 url: str, 

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

1314 pool_manager: PoolManager | None = None, 

1315 ) -> HTTPResponse: 

1316 """Send a webDAV MKCOL request and return the response unmodified. 

1317 

1318 Parameters 

1319 ---------- 

1320 url : `str` 

1321 Target URL. 

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

1323 Headers to sent with the request. 

1324 pool_manager : `PoolManager`, optional 

1325 Pool manager to use to send this request. 

1326 

1327 Notes 

1328 ----- 

1329 This method is intended for subclasses to override when needed. 

1330 """ 

1331 return self._request("MKCOL", url=url, headers=headers, pool_manager=pool_manager) 

1332 

1333 def _move( 

1334 self, 

1335 url: str, 

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

1337 pool_manager: PoolManager | None = None, 

1338 ) -> HTTPResponse: 

1339 """Send a webDAV MOVE request and return the response unmodified. 

1340 

1341 Parameters 

1342 ---------- 

1343 url : `str` 

1344 Target URL. 

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

1346 Headers to sent with the request. 

1347 pool_manager : `PoolManager`, optional 

1348 Pool manager to use to send this request. 

1349 

1350 Notes 

1351 ----- 

1352 This method is intended for subclasses to override when needed. 

1353 """ 

1354 return self._request("MOVE", url=url, headers=headers, pool_manager=pool_manager) 

1355 

1356 def _propfind( 

1357 self, 

1358 url: str, 

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

1360 body: str = "", 

1361 pool_manager: PoolManager | None = None, 

1362 ) -> HTTPResponse: 

1363 """Send a webDAV PROPFIND request and return the response unmodified. 

1364 

1365 Parameters 

1366 ---------- 

1367 url : `str` 

1368 Target URL. 

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

1370 Headers to sent with the request. 

1371 body : `str`, optional 

1372 Request body. 

1373 pool_manager : `PoolManager`, optional 

1374 Pool manager to use to send this request. 

1375 

1376 Notes 

1377 ----- 

1378 This method is intended for subclasses to override when needed. 

1379 """ 

1380 return self._request("PROPFIND", url=url, headers=headers, body=body, pool_manager=pool_manager) 

1381 

1382 def _put( 

1383 self, 

1384 url: str, 

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

1386 body: BinaryIO | bytes = b"", 

1387 preload_content: bool = True, 

1388 redirect: bool = True, 

1389 pool_manager: PoolManager | None = None, 

1390 ) -> HTTPResponse: 

1391 """Send a HTTP PUT request and return the response unmodified. 

1392 

1393 Parameters 

1394 ---------- 

1395 url : `str` 

1396 Target URL. 

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

1398 Headers to sent with the request. 

1399 body : `BinaryIO` or `bytes`, optional 

1400 Request body. 

1401 preload_content : `bool`, optional 

1402 If True, the response body is downloaded and can be retrieved 

1403 via the returned response `.data` property. If False, the 

1404 caller needs to call the `.read()` on the returned response 

1405 object to download the body. 

1406 redirect : `bool`, optional 

1407 If True, follow redirections. 

1408 pool_manager : `PoolManager`, optional 

1409 Pool manager to send the request through. 

1410 

1411 Returns 

1412 ------- 

1413 resp: `HTTPResponse` 

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

1415 

1416 Notes 

1417 ----- 

1418 This method is intended for subclasses to override when needed. 

1419 """ 

1420 # Disable retries when we know the request is not idempotent. In 

1421 # particular, when the body of the request is an `io.BufferedReader`, 

1422 # any attempt to use that body may totally or partially consume it. 

1423 # That means that in case of a retry, the last successful attempt may 

1424 # end up uploading an incomplete body and, as a consequence, the 

1425 # resulting uploaded file may be either incomplete or have a length 

1426 # of zero. 

1427 # 

1428 # So we only retry a PUT request when the body is an instance of 

1429 # `bytes`. 

1430 # 

1431 # Note that we cannot set `retries` to False since that setting would 

1432 # also disable redirection. To disable retries only, we must explicitly 

1433 # set the `retries` keyword argument to integer value zero (0). 

1434 # 

1435 # See documentation: 

1436 # https://urllib3.readthedocs.io/en/stable/user-guide.html#retrying-requests 

1437 kwargs: dict[Any, Any] = {} if isinstance(body, bytes) else {"retries": 0} 

1438 

1439 return self._request( 

1440 "PUT", 

1441 url=url, 

1442 headers=headers, 

1443 body=body, 

1444 preload_content=preload_content, 

1445 redirect=redirect, 

1446 pool_manager=pool_manager, 

1447 **kwargs, 

1448 ) 

1449 

1450 def head( 

1451 self, 

1452 url: str, 

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

1454 ) -> HTTPResponse: 

1455 """Send a HTTP HEAD request, process and return the response 

1456 only if successful. 

1457 

1458 Parameters 

1459 ---------- 

1460 url : `str` 

1461 Target URL. 

1462 headers : `bool` 

1463 If the target URL is not found, raise an exception. Otherwise 

1464 just return the response. 

1465 """ 

1466 headers = {} if headers is None else dict(headers) 

1467 resp = self._head(url=url, headers=headers) 

1468 match resp.status: 

1469 case HTTPStatus.OK: 

1470 return resp 

1471 case HTTPStatus.NOT_FOUND: 

1472 raise FileNotFoundError(f"No file found at {resp.geturl()}") 

1473 case _: 

1474 raise unexpected_status_error("HEAD", url, resp) 

1475 

1476 def get( 

1477 self, 

1478 url: str, 

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

1480 preload_content: bool = True, 

1481 redirect: bool = True, 

1482 ) -> tuple[str, HTTPResponse]: 

1483 """Send a HTTP GET request. 

1484 

1485 Parameters 

1486 ---------- 

1487 url : `str` 

1488 Target URL. 

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

1490 Headers to sent with the request. 

1491 preload_content : `bool`, optional 

1492 If True, the response body is downloaded and can be retrieved 

1493 via the returned response `.data` property. If False, the 

1494 caller needs to call the `.read()` on the returned response 

1495 object to download the body. 

1496 redirect : `bool`, optional 

1497 If True, follow redirections. 

1498 

1499 Returns 

1500 ------- 

1501 url: `str` 

1502 The URL we used to obtain this response. It may be different from 

1503 the URL passed as argument in case of redirection. 

1504 resp: `HTTPResponse` 

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

1506 """ 

1507 # Send the GET request to the frontend servers. 

1508 headers = {} if headers is None else dict(headers) 

1509 resp = self._get( 

1510 url, 

1511 headers=headers, 

1512 preload_content=preload_content, 

1513 redirect=redirect, 

1514 ) 

1515 match resp.status: 

1516 case HTTPStatus.OK | HTTPStatus.PARTIAL_CONTENT: 

1517 return self._get_response_url(resp, default_url=url), resp 

1518 case HTTPStatus.NOT_FOUND: 

1519 raise FileNotFoundError(f"No file found at {resp.geturl()}") 

1520 case status if status in resp.REDIRECT_STATUSES and not redirect: 

1521 # This response is a redirection but we are asked not to 

1522 # follow redirections, so return this response as is. 

1523 return self._get_response_url(resp, default_url=url), resp 

1524 case _: 

1525 raise unexpected_status_error("GET", url, resp) 

1526 

1527 def options( 

1528 self, 

1529 url: str, 

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

1531 ) -> HTTPResponse: 

1532 """Send a HTTP OPTIONS request and return the response on success. 

1533 

1534 Parameters 

1535 ---------- 

1536 url : `str` 

1537 Target URL. 

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

1539 Headers to sent with the request. 

1540 

1541 Returns 

1542 ------- 

1543 resp: `HTTPResponse` 

1544 Response to the request as received from the server. 

1545 """ 

1546 resp = self._options(url=url, headers=headers) 

1547 match resp.status: 

1548 case HTTPStatus.OK | HTTPStatus.CREATED: 

1549 return resp 

1550 case _: 

1551 raise unexpected_status_error("OPTIONS", url, resp) 

1552 

1553 def propfind( 

1554 self, 

1555 url: str, 

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

1557 body: str = "", 

1558 depth: str = "0", 

1559 ) -> HTTPResponse: 

1560 """Send a HTTP PROPFIND request and return the unmodified response on 

1561 success. 

1562 

1563 Parameters 

1564 ---------- 

1565 url : `str` 

1566 Target URL. 

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

1568 Headers to sent with the request. 

1569 body : `str`, optional 

1570 Request body. 

1571 depth : `str`, optional 

1572 ???. 

1573 """ 

1574 headers = {} if headers is None else dict(headers) 

1575 headers.update( 

1576 { 

1577 "Depth": depth, 

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

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

1580 } 

1581 ) 

1582 resp = self._propfind(url=url, headers=headers, body=body) 

1583 match resp.status: 

1584 case HTTPStatus.MULTI_STATUS | HTTPStatus.NOT_FOUND: 

1585 return resp 

1586 case _: 

1587 raise unexpected_status_error("PROPFIND", url, resp) 

1588 

1589 def put( 

1590 self, 

1591 url: str, 

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

1593 data: BinaryIO | bytes = b"", 

1594 ) -> int | None: 

1595 """Send a HTTP PUT request. 

1596 

1597 Parameters 

1598 ---------- 

1599 url : `str` 

1600 Target URL. 

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

1602 Headers to sent with the request. 

1603 data : `BinaryIO` or `bytes` 

1604 Request body. 

1605 

1606 Returns 

1607 ------- 

1608 size : `int | None` 

1609 The size in bytes of the file uploaded. Can be `None` if the size 

1610 could not be retrieved. 

1611 """ 

1612 # Send a PUT request with empty body and handle redirection. This 

1613 # is useful if the server redirects us; since we cannot rewind the 

1614 # data we are uploading, we don't start uploading data until we 

1615 # connect to the server that will actually serve our request. 

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

1617 frontend_headers.update({"Content-Length": "0"}) 

1618 resp = self._put(url, headers=frontend_headers, body=b"", redirect=False) 

1619 match resp.status: 

1620 case HTTPStatus.OK | HTTPStatus.CREATED | HTTPStatus.NO_CONTENT: 

1621 redirect_url = url 

1622 case status if status in resp.REDIRECT_STATUSES: 

1623 redirect_url = resp.headers.get("Location") 

1624 case _: 

1625 raise unexpected_status_error("PUT", url, resp) 

1626 

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

1628 # its final destination. 

1629 

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

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

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

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

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

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

1636 # 

1637 # See RFC-3230 for details and 

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

1639 # for the list of supported digest algorithhms. 

1640 # 

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

1642 # the checksum reqquest may be ignored by the server. 

1643 backend_headers = {} if headers is None else dict(headers) 

1644 if (checksum := self._config.request_checksum) is not None: 

1645 backend_headers.update({"Want-Digest": checksum}) 

1646 

1647 resp = self._put(redirect_url, body=data, headers=backend_headers) 

1648 match resp.status: 

1649 case HTTPStatus.OK | HTTPStatus.CREATED | HTTPStatus.NO_CONTENT: 

1650 # Send a HEAD request to retrieve the size of the file we 

1651 # just uploaded 

1652 resp = self.head(redirect_url) 

1653 size = int(resp.headers.get("Content-Length", -1)) 

1654 return None if size == -1 else size 

1655 case _: 

1656 raise unexpected_status_error("PUT", redirect_url, resp) 

1657 

1658 def _get_temporary_basename(self, basename: str, prefix: str) -> str: 

1659 """Return a basename for a temporary file.""" 

1660 unique_id = str(uuid.uuid4()) 

1661 return f"{prefix}.{unique_id}.{basename}" 

1662 

1663 def _split_parent_and_basename(self, url: str) -> tuple[str, str]: 

1664 """Return the URL of the parent directory and the basename from 

1665 `url`. 

1666 """ 

1667 parsed: Url = parse_url(url) 

1668 normalized_path = normalize_path(parsed.path) 

1669 parent_path = posixpath.dirname(normalized_path) 

1670 basename = posixpath.basename(normalized_path) 

1671 parent_url = Url( 

1672 scheme=parsed.scheme, 

1673 auth=parsed.auth, 

1674 host=parsed.host, 

1675 port=parsed.port, 

1676 path=parent_path, 

1677 query=parsed.query, 

1678 fragment=parsed.fragment, 

1679 ).url 

1680 return parent_url, basename 

1681 

1682 def _parent(self, url: str) -> str: 

1683 """Return the URL of the parent directory to `url`.""" 

1684 parent_url, _ = self._split_parent_and_basename(url) 

1685 return parent_url 

1686 

1687 def _make_temporary_url(self, url: str, prefix: str = ".tmp") -> str: 

1688 """Return the URL of a temporary file based on `url`.""" 

1689 parent_url, basename = self._split_parent_and_basename(url) 

1690 temporary_basename = self._get_temporary_basename(basename=basename, prefix=prefix) 

1691 return f"{parent_url}/{temporary_basename}" 

1692 

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

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

1695 

1696 Parameters 

1697 ---------- 

1698 url : `str` 

1699 Target URL. 

1700 

1701 Returns 

1702 ------- 

1703 result: `bool` 

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

1705 """ 

1706 return self.stat(url).exists 

1707 

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

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

1710 

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

1712 

1713 Parameters 

1714 ---------- 

1715 url : `str` 

1716 Target URL. 

1717 

1718 Returns 

1719 ------- 

1720 size: `int` 

1721 The number of bytes of the resource located at `url`. 

1722 """ 

1723 # Check if we have the size of this URL in our cache 

1724 if (size := self._file_size_cache.get_size(url)) is not None: 

1725 return size 

1726 

1727 stat = self.stat(url) 

1728 if not stat.exists: 

1729 raise FileNotFoundError(f"No file or directory found at {url}") 

1730 else: 

1731 return stat.size 

1732 

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

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

1735 

1736 Parameters 

1737 ---------- 

1738 url : `str` 

1739 Target URL. 

1740 

1741 Returns 

1742 ------- 

1743 result: `bool` 

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

1745 """ 

1746 return self.stat(url).is_dir 

1747 

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

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

1750 

1751 If a directory already exists at `url` no error is returned nor 

1752 exception is raised. An exception is raised if a file exists at `url`. 

1753 

1754 Parameters 

1755 ---------- 

1756 url : `str` 

1757 Target URL. 

1758 """ 

1759 resp = self._mkcol(url=url) 

1760 match resp.status: 

1761 case HTTPStatus.CREATED | HTTPStatus.METHOD_NOT_ALLOWED: 

1762 return 

1763 case HTTPStatus.CONFLICT: 

1764 # The parent directory does not exist. Create it first except 

1765 # if the parent's path is "/". 

1766 parent = self._parent(url) 

1767 if not parent.endswith("/"): 

1768 self.mkcol(parent) 

1769 resp = self._mkcol(url=url) 

1770 case _: 

1771 raise ValueError( 

1772 f"Can not create directory {resp.geturl()}: status {resp.status} {resp.reason}" 

1773 ) 

1774 

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

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

1777 

1778 Parameters 

1779 ---------- 

1780 url : `str` 

1781 Target URL. 

1782 

1783 Returns 

1784 ------- 

1785 result: `DavResourceMetadata` 

1786 Details of the resources at `url`. If no resource was found at 

1787 that URL no exception is raised. Instead the returned details allow 

1788 for detecting that the resource does not exist. 

1789 

1790 The returned value should include fields to determine 

1791 if there is a file or a directory at that `url` and if so, its 

1792 size and kind (file or directory). Other fields may also be 

1793 included depending on the implementation of the webDAV protocol 

1794 by the server. 

1795 """ 

1796 # Request the minimum set of DAV properties. 

1797 body = ( 

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

1799 """<D:propfind xmlns:D="DAV:">""" 

1800 """<D:prop>""" 

1801 """<D:resourcetype/>""" 

1802 """<D:getcontentlength/>""" 

1803 """<D:getlastmodified/>""" 

1804 """</D:prop>""" 

1805 """</D:propfind>""" 

1806 ) 

1807 resp = self.propfind(url, body=body, depth="0") 

1808 match resp.status: 

1809 case HTTPStatus.NOT_FOUND: 

1810 href = url.replace(self._base_url, "", 1) 

1811 return DavFileMetadata(base_url=self._base_url, href=href) 

1812 case HTTPStatus.MULTI_STATUS: 

1813 property = self._propfind_parser.parse(resp.data)[0] 

1814 return DavFileMetadata.from_property(base_url=self._base_url, property=property) 

1815 case _: 

1816 raise unexpected_status_error("PROPFIND", url, resp) 

1817 

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

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

1820 

1821 Parameters 

1822 ---------- 

1823 url : `str` 

1824 Target URL. 

1825 name : `str` 

1826 Name of the object to be included in the returned value. If None, 

1827 the `url` is used as name. 

1828 

1829 Returns 

1830 ------- 

1831 result: `dict` 

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

1833 

1834 .. code-block:: json 

1835 

1836 { 

1837 "name": name, 

1838 "size": 1234, 

1839 "type": "file", 

1840 "last_modified": 

1841 datetime.datetime(2025, 4, 10, 15, 12, 51, 227854), 

1842 "checksums": { 

1843 "adler32": "0fc5f83f", 

1844 "md5": "1f57339acdec099c6c0a41f8e3d5fcd0", 

1845 } 

1846 } 

1847 

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

1849 

1850 .. code-block:: json 

1851 

1852 { 

1853 "name": name, 

1854 "size": 0, 

1855 "type": "directory", 

1856 "last_modified": 

1857 datetime.datetime(2025, 4, 10, 15, 12, 51, 227854), 

1858 "checksums": {}, 

1859 } 

1860 

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

1862 form: 

1863 

1864 .. code-block:: json 

1865 

1866 { 

1867 "name": name, 

1868 "size": None, 

1869 "type": None, 

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

1871 "checksums": {}, 

1872 } 

1873 

1874 Notes 

1875 ----- 

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

1877 `fsspec`. 

1878 

1879 The size of existing directories is always zero. The `checksums` 

1880 dictionary is empty for directories and may be empty for files if the 

1881 server does not compute and store the checksum of the files it stores. 

1882 """ 

1883 result: dict[str, Any] = { 

1884 "name": name if name is not None else url, 

1885 "type": None, 

1886 "size": None, 

1887 "last_modified": datetime.min, 

1888 "checksums": {}, 

1889 } 

1890 metadata = self.stat(url) 

1891 if not metadata.exists: 

1892 return result 

1893 

1894 result.update( 

1895 { 

1896 "type": "directory" if metadata.is_dir else "file", 

1897 "size": metadata.size, 

1898 "last_modified": metadata.last_modified, 

1899 "checksums": metadata.checksums, 

1900 } 

1901 ) 

1902 return result 

1903 

1904 def move(self, source_url: str, destination_url: str, overwrite: bool = False) -> HTTPResponse: 

1905 """Send a webDAV MOVE request and return the response unmodified. 

1906 

1907 Parameters 

1908 ---------- 

1909 source_url : `str` 

1910 Source URL. 

1911 destination_url : `str` 

1912 Destination URL. 

1913 overwrite : `bool`, optional 

1914 Overwrite the destination if it exists. 

1915 

1916 Returns 

1917 ------- 

1918 resp : `HTTPResponse` 

1919 The unmodified response received from the server. 

1920 """ 

1921 headers = { 

1922 "Destination": destination_url, 

1923 "Overwrite": "T" if overwrite else "F", 

1924 } 

1925 return self._move(source_url, headers=headers) 

1926 

1927 def read_dir(self, url: str) -> list[DavFileMetadata]: 

1928 """Return the properties of the files or directories contained in 

1929 directory located at `url`. 

1930 

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

1932 

1933 Parameters 

1934 ---------- 

1935 url : `str` 

1936 Target URL. 

1937 

1938 Returns 

1939 ------- 

1940 result: `list[DavResourceMetadata]` 

1941 List of details of each file or directory within `url`. 

1942 """ 

1943 body = ( 

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

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

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

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

1948 ) 

1949 resp = self.propfind(url, body=body, depth="1") 

1950 match resp.status: 

1951 case HTTPStatus.MULTI_STATUS: 

1952 pass 

1953 case HTTPStatus.NOT_FOUND: 

1954 raise FileNotFoundError(f"No directory found at {resp.geturl()}") 

1955 case _: 

1956 raise unexpected_status_error("PROPFIND", url, resp) 

1957 

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

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

1960 else: 

1961 this_dir_href = "/" 

1962 

1963 result = [] 

1964 for property in self._propfind_parser.parse(resp.data): 

1965 # Don't include in the results the metadata of the directory we 

1966 # traversing. 

1967 # Some webDAV servers do not append a "/" to the href of a 

1968 # directory in their response to PROPFIND, so we must take into 

1969 # account that. 

1970 if property.is_file: 

1971 result.append(DavFileMetadata.from_property(base_url=self._base_url, property=property)) 

1972 elif property.is_dir and property.href != this_dir_href: 

1973 result.append(DavFileMetadata.from_property(base_url=self._base_url, property=property)) 

1974 

1975 return result 

1976 

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

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

1979 

1980 Parameters 

1981 ---------- 

1982 url : `str` 

1983 Target URL. 

1984 

1985 Returns 

1986 ------- 

1987 url: `str` 

1988 Backend URL from which the data was obtained. 

1989 data: `bytes` 

1990 Contents of the file. 

1991 

1992 Notes 

1993 ----- 

1994 The caller must ensure that the resource at `url` is a file, not 

1995 a directory. 

1996 """ 

1997 backend_url, resp = self.get(url) 

1998 return backend_url, resp.data 

1999 

2000 def read_range( 

2001 self, 

2002 url: str, 

2003 start: int, 

2004 end: int | None, 

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

2006 ) -> tuple[str, bytes]: 

2007 """Download partial content of file located at `url`. 

2008 

2009 Parameters 

2010 ---------- 

2011 url : `str` 

2012 Target URL. 

2013 start : `int` 

2014 Starting byte offset of the range to download. 

2015 end : `int`, optional 

2016 Ending byte offset of the range to download. 

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

2018 Specific headers to sent with the GET request. 

2019 

2020 Returns 

2021 ------- 

2022 backend_url: `str` 

2023 URL used to retrieve this data. If the server redirected us 

2024 this is the URL we were redirected to. 

2025 data: `bytes` 

2026 Partial contents of the file. 

2027 

2028 Notes 

2029 ----- 

2030 The caller must ensure that the resource at `url` is a file, not 

2031 a directory. This is important because some webDAV servers respond 

2032 with an HTML document when asked for reading a directory. 

2033 """ 

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

2035 get_headers.update({"Accept-Encoding": "identity"}) 

2036 if end is None: 

2037 get_headers.update({"Range": f"bytes={start}-"}) 

2038 else: 

2039 get_headers.update({"Range": f"bytes={start}-{end}"}) 

2040 

2041 final_url, resp = self.get(url, headers=get_headers, redirect=True) 

2042 match resp.status: 

2043 case HTTPStatus.PARTIAL_CONTENT: 

2044 return final_url, resp.data 

2045 case _: 

2046 raise unexpected_status_error("GET (with 'Range' header)", url, resp) 

2047 

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

2049 """Signal the server hosting `url` that no more `read_range` requests 

2050 will be issued, so that the resources allocated for serving those 

2051 requests can be released. 

2052 

2053 This helper method is intended to be overwritten by subclasses that 

2054 need such functionality according to the behavior of the specific 

2055 remote storage endpoint. 

2056 

2057 Parameters 

2058 ---------- 

2059 url : `str` 

2060 Target URL. 

2061 """ 

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

2063 pass 

2064 

2065 def _write_response_body_to_file(self, resp: HTTPResponse, filename: str, chunk_size: int) -> int: 

2066 """Write the response body to a local file. 

2067 

2068 Parameters 

2069 ---------- 

2070 resp : `HTTPResponse` 

2071 The HTTP Response to read the body from. 

2072 filename : `str` 

2073 Local file to write the content to. If the file already exists, 

2074 it will be rewritten. 

2075 chunk_size : `int` 

2076 Size of the chunks to write to `filename`. 

2077 

2078 Returns 

2079 ------- 

2080 count: `int` 

2081 Number of bytes written to `filename`. 

2082 """ 

2083 try: 

2084 # Read the response body into a pre-allocated memory buffer and 

2085 # write the buffer content to the destination file avoiding 

2086 # copies if possible. 

2087 content_length = 0 

2088 with open(filename, "wb", buffering=0) as file: 

2089 view = memoryview(bytearray(chunk_size)) 

2090 while True: 

2091 if (count := resp.readinto(view)) > 0: # type: ignore 

2092 content_length += count 

2093 file.write(view[:count]) 

2094 else: 

2095 break 

2096 

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

2098 # Perform this check only when the body of the response was not 

2099 # encoded by the server. 

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

2101 if ( 

2102 "Content-Encoding" not in resp.headers 

2103 and expected_length != -1 

2104 and expected_length != content_length 

2105 ): 

2106 raise ValueError( 

2107 f"Size of downloaded file does not match value in Content-Length header for " 

2108 f"{resp.geturl()}: expecting {expected_length} and got {content_length} bytes" 

2109 ) 

2110 

2111 return content_length 

2112 finally: 

2113 # Release the connection 

2114 resp.drain_conn() 

2115 resp.release_conn() 

2116 

2117 def download(self, url: str, filename: str, chunk_size: int) -> int: 

2118 """Download the content of a file and write it to local file. 

2119 

2120 Parameters 

2121 ---------- 

2122 url : `str` 

2123 Target URL. 

2124 filename : `str` 

2125 Local file to write the content to. If the file already exists, 

2126 it will be rewritten. 

2127 chunk_size : `int` 

2128 Size of the chunks to write to `filename`. 

2129 

2130 Returns 

2131 ------- 

2132 count: `int` 

2133 Number of bytes written to `filename`. 

2134 

2135 Notes 

2136 ----- 

2137 The caller must ensure that the resource at `url` is a file, not 

2138 a directory. 

2139 """ 

2140 _, resp = self.get(url, preload_content=False) 

2141 return self._write_response_body_to_file(resp, filename, chunk_size) 

2142 

2143 def write(self, url: str, data: BinaryIO | bytes) -> int | None: 

2144 """Create or rewrite a remote file at `url` with `data` as its 

2145 contents. 

2146 

2147 Parameters 

2148 ---------- 

2149 url : `str` 

2150 Target URL. 

2151 data : `bytes` 

2152 Sequence of bytes to upload. 

2153 

2154 Returns 

2155 ------- 

2156 size : `int | None` 

2157 The size in bytes of the file uploaded. Can be `None` if the size 

2158 could not be retrieved. 

2159 

2160 Notes 

2161 ----- 

2162 If a file already exists at `url` it will be rewritten. 

2163 """ 

2164 # According to RFC 4918, the parent directory of the file must 

2165 # exist before we can write to it. So create it first and then 

2166 # upload. 

2167 self.mkcol(self._parent(url)) 

2168 

2169 try: 

2170 # Upload to a temporary file and rename to the final name. 

2171 temporary_url = self._make_temporary_url(url) 

2172 size = self.put(temporary_url, data=data) 

2173 self.rename(temporary_url, url, overwrite=True, create_parent=False) 

2174 

2175 # Update the file size cache with this size 

2176 self._file_size_cache.update_size(url, size) 

2177 return size 

2178 except Exception: 

2179 # Upload failed. Attempt to remove the temporary file. 

2180 self.delete(temporary_url) 

2181 raise 

2182 

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

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

2185 

2186 The checksums are retrieved from the storage endpoint. There may be 

2187 none if the storage endpoint does not automatically expose the 

2188 checksums it computes. 

2189 

2190 Parameters 

2191 ---------- 

2192 url : `str` 

2193 Target URL. 

2194 

2195 Returns 

2196 ------- 

2197 checksums: `dict[str, str]` 

2198 A file exists at `url`. 

2199 The key of the dictionary is the lowercased name of the checksum 

2200 algorithm (e.g. "md5", "adler32"). The value is the lowercased 

2201 checksum itself (e.g. "78441cec2479ec8b545c4d6699f542da"). 

2202 """ 

2203 stat = self.stat(url) 

2204 if not stat.exists: 

2205 raise FileNotFoundError(f"No file found at {url}") 

2206 

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

2208 

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

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

2211 

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

2213 

2214 Parameters 

2215 ---------- 

2216 url : `str` 

2217 Target URL. 

2218 

2219 Notes 

2220 ----- 

2221 If `url` designates a directory, some webDAV servers recursively 

2222 remove the directory and its contents. Others, only remove the 

2223 directory if it is empty. 

2224 

2225 For a consisten behavior, the caller must check what kind of object 

2226 the target URL is and walk the hierarchy removing all objects. 

2227 """ 

2228 resp = self._delete(url) 

2229 match resp.status: 

2230 case HTTPStatus.OK | HTTPStatus.ACCEPTED | HTTPStatus.NO_CONTENT | HTTPStatus.NOT_FOUND: 

2231 # Invalidate the entry for this file in our cache, if any 

2232 self._file_size_cache.invalidate(url) 

2233 case _: 

2234 raise ValueError( 

2235 f"Unable to delete resource {resp.geturl()}: status {resp.status} {resp.reason}" 

2236 ) 

2237 

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

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

2240 GET requests against `url`. 

2241 

2242 Parameters 

2243 ---------- 

2244 url : `str` 

2245 Target URL. 

2246 """ 

2247 # If we have already determined that the server accepts "Range" for 

2248 # another URL, we assume that it implements that feature for any 

2249 # file it serves, so reuse that information. 

2250 if self._accepts_ranges is not None: 

2251 return self._accepts_ranges 

2252 

2253 with self._lock: 

2254 if self._accepts_ranges is None: 

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

2256 

2257 return self._accepts_ranges 

2258 

2259 @property 

2260 def supports_duplicate(self) -> bool: 

2261 """Return True if the server this client interacts with implements 

2262 webDAV COPY method. 

2263 """ 

2264 return self._can_duplicate 

2265 

2266 def copy(self, source_url: str, destination_url: str, overwrite: bool = False) -> None: 

2267 """Copy the file at `source_url` to `destination_url` in the same 

2268 storage endpoint. 

2269 

2270 Parameters 

2271 ---------- 

2272 source_url : `str` 

2273 URL of the source file. 

2274 destination_url : `str` 

2275 URL of the destination file. Its parent directory must exist. 

2276 overwrite : `bool` 

2277 If True and a file exists at `destination_url` it will be 

2278 overwritten. Otherwise an exception is raised. 

2279 """ 

2280 headers = { 

2281 "Destination": destination_url, 

2282 "Overwrite": "T" if overwrite else "F", 

2283 } 

2284 resp = self._copy(source_url, headers=headers) 

2285 match resp.status: 

2286 case HTTPStatus.CREATED | HTTPStatus.NO_CONTENT: 

2287 self._file_size_cache.invalidate(destination_url) 

2288 case _: 

2289 raise ValueError( 

2290 f"Could not copy {resp.geturl()} to {destination_url}: status {resp.status} {resp.reason}" 

2291 ) 

2292 

2293 def duplicate(self, source_url: str, destination_url: str, overwrite: bool = False) -> None: 

2294 """Copy the file at `source_url` to `destination_url` in the same 

2295 storage endpoint. 

2296 

2297 Parameters 

2298 ---------- 

2299 source_url : `str` 

2300 URL of the source file. 

2301 destination_url : `str` 

2302 URL of the destination file. Its parent directory is created if 

2303 necessary. 

2304 overwrite : `bool` 

2305 If True and a file exists at `destination_url` it will be 

2306 overwritten. Otherwise an exception is raised. 

2307 """ 

2308 # Check the source is a file 

2309 if self.is_dir(source_url): 

2310 raise NotImplementedError(f"copy is not implemented for directory {source_url}") 

2311 

2312 # Create the destination's parent directory first because COPY may 

2313 # fail if it does not exist, depending on the server implementation 

2314 # of RFC 4918. 

2315 destination_parent = self._parent(destination_url) 

2316 self.mkcol(destination_parent) 

2317 self.copy(source_url=source_url, destination_url=destination_url, overwrite=overwrite) 

2318 

2319 def rename( 

2320 self, 

2321 source_url: str, 

2322 destination_url: str, 

2323 overwrite: bool = False, 

2324 create_parent: bool = True, 

2325 ) -> None: 

2326 """Rename (move) the file at `source_url` to `destination_url` in the 

2327 same storage endpoint. 

2328 

2329 Parameters 

2330 ---------- 

2331 source_url : `str` 

2332 URL of the source file. 

2333 destination_url : `str` 

2334 URL of the destination file. Its parent directory must exist. 

2335 overwrite : `bool`, optional 

2336 If True and a file exists at `destination_url` it will be 

2337 overwritten. Otherwise an exception is raised. 

2338 create_parent : `bool`, optional 

2339 Whether to create the parent. 

2340 """ 

2341 # Create the destination's parent directory first because MOVE may 

2342 # fail if it does not exist, depending on the server implementation 

2343 # of RFC 4918. 

2344 if create_parent: 

2345 destination_parent = self._parent(destination_url) 

2346 self.mkcol(destination_parent) 

2347 

2348 resp = self.move(source_url=source_url, destination_url=destination_url, overwrite=overwrite) 

2349 match resp.status: 

2350 case HTTPStatus.OK | HTTPStatus.CREATED | HTTPStatus.NO_CONTENT: 

2351 self._file_size_cache.invalidate(destination_url) 

2352 case _: 

2353 raise ValueError( 

2354 f"""Could not move file {resp.geturl()} to {destination_url}: status {resp.status} """ 

2355 f"""{resp.reason}""" 

2356 ) 

2357 

2358 def generate_presigned_get_url(self, url: str, expiration_time_seconds: int) -> str: 

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

2360 using an HTTP GET without supplying any access credentials. 

2361 

2362 Parameters 

2363 ---------- 

2364 url : `str` 

2365 Target URL. 

2366 expiration_time_seconds : `int` 

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

2368 

2369 Returns 

2370 ------- 

2371 url : `str` 

2372 HTTP URL signed for GET. 

2373 """ 

2374 raise NotImplementedError(f"URL signing is not supported by server for {self}") 

2375 

2376 def generate_presigned_put_url(self, url: str, expiration_time_seconds: int) -> str: 

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

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

2379 

2380 Parameters 

2381 ---------- 

2382 url : `str` 

2383 Target URL. 

2384 expiration_time_seconds : `int` 

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

2386 

2387 Returns 

2388 ------- 

2389 url : `str` 

2390 HTTP URL signed for PUT. 

2391 """ 

2392 raise NotImplementedError(f"URL signing is not supported by server for {self}") 

2393 

2394 

2395class ActivityCaveat(enum.Enum): 

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

2397 macaroons for dCache or XRootD webDAV servers. 

2398 """ 

2399 

2400 DOWNLOAD = 1 

2401 UPLOAD = 2 

2402 

2403 

2404class DavClientURLSigner(DavClient): 

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

2406 

2407 Instances of this class are thread-safe. 

2408 

2409 Parameters 

2410 ---------- 

2411 url : `str` 

2412 Root URL of the storage endpoint 

2413 (e.g. "https://host.example.org:1234/"). 

2414 config : `DavConfig` 

2415 Configuration to initialize this client. 

2416 accepts_ranges : `bool` | `None` 

2417 Indicate whether the remote server accepts the ``Range`` header in GET 

2418 requests. 

2419 """ 

2420 

2421 def __init__(self, url: str, config: DavConfig, accepts_ranges: bool | None = None) -> None: 

2422 super().__init__(url=url, config=config, accepts_ranges=accepts_ranges) 

2423 

2424 def generate_presigned_get_url(self, url: str, expiration_time_seconds: int) -> str: 

2425 """Return a pre-signed URL that can be used to retrieve the resource 

2426 at `url` using an HTTP GET without supplying any access credentials. 

2427 

2428 Parameters 

2429 ---------- 

2430 url : `str` 

2431 URL of an existing file. 

2432 expiration_time_seconds : `int` 

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

2434 

2435 Returns 

2436 ------- 

2437 url : `str` 

2438 HTTP URL signed for GET. 

2439 

2440 Notes 

2441 ----- 

2442 Although the returned URL allows for downloading the file at `url` 

2443 without supplying credentials, the HTTP client must be configured 

2444 to accept the certificate the server will present if the client wants 

2445 validate it. The server's certificate may be issued by a certificate 

2446 authority unknown to the client. 

2447 """ 

2448 macaroon: str = self._get_macaroon(url, ActivityCaveat.DOWNLOAD, expiration_time_seconds) 

2449 return f"{url}?authz={macaroon}" 

2450 

2451 def generate_presigned_put_url(self, url: str, expiration_time_seconds: int) -> str: 

2452 """Return a pre-signed URL that can be used to upload a file to `url` 

2453 using an HTTP PUT without supplying any access credentials. 

2454 

2455 Parameters 

2456 ---------- 

2457 url : `str` 

2458 URL of an existing file. 

2459 expiration_time_seconds : `int` 

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

2461 

2462 Returns 

2463 ------- 

2464 url : `str` 

2465 HTTP URL signed for PUT. 

2466 

2467 Notes 

2468 ----- 

2469 Although the returned URL allows for uploading a file to `url` 

2470 without supplying credentials, the HTTP client must be configured 

2471 to accept the certificate the server will present if the client wants 

2472 validate it. The server's certificate may be issued by a certificate 

2473 authority unknown to the client. 

2474 """ 

2475 macaroon: str = self._get_macaroon(url, ActivityCaveat.UPLOAD, expiration_time_seconds) 

2476 return f"{url}?authz={macaroon}" 

2477 

2478 def _get_macaroon(self, url: str, activity: ActivityCaveat, expiration_time_seconds: int) -> str: 

2479 """Return a macaroon for uploading or downloading the file at `url`. 

2480 

2481 Parameters 

2482 ---------- 

2483 url : `str` 

2484 URL of an existing file. 

2485 activity : `ActivityCaveat` 

2486 the activity the macaroon is requested for. 

2487 expiration_time_seconds : `int` 

2488 Requested duration of the macaroon, in seconds. 

2489 

2490 Returns 

2491 ------- 

2492 macaroon : `str` 

2493 Macaroon to be used with `url` in a GET or PUT request. 

2494 """ 

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

2496 # 

2497 # For details about dCache macaroons see: 

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

2499 match activity: 

2500 case ActivityCaveat.DOWNLOAD: 

2501 activity_caveat = "DOWNLOAD,LIST" 

2502 case ActivityCaveat.UPLOAD: 

2503 activity_caveat = "UPLOAD,LIST,DELETE,MANAGE" 

2504 

2505 # Retrieve a macaroon for the requested activities and duration 

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

2507 body = { 

2508 "caveats": [ 

2509 f"activity:{activity_caveat}", 

2510 ], 

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

2512 } 

2513 resp = self._request("POST", url, headers=headers, body=json.dumps(body)) 

2514 if resp.status != HTTPStatus.OK: 

2515 raise ValueError( 

2516 f"Could not retrieve a macaroon for URL {resp.geturl()}, status: {resp.status} {resp.reason}" 

2517 ) 

2518 

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

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

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

2522 # 

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

2524 # { 

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

2526 # "uri": { 

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

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

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

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

2531 # } 

2532 # } 

2533 # 

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

2535 # { 

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

2537 # "expires_in": 86400 

2538 # } 

2539 try: 

2540 response_body = json.loads(resp.data.decode()) 

2541 except json.JSONDecodeError: 

2542 raise ValueError(f"Could not deserialize response to POST request for URL {resp.geturl()}") 

2543 

2544 if "macaroon" in response_body: 

2545 return response_body["macaroon"] 

2546 

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

2548 

2549 @override 

2550 def duplicate(self, source_url: str, destination_url: str, overwrite: bool = False) -> None: 

2551 """Copy the file at `source_url` to `destination_url` in the same 

2552 storage endpoint. 

2553 

2554 Parameters 

2555 ---------- 

2556 source_url : `str` 

2557 URL of the source file. 

2558 destination_url : `str` 

2559 URL of the destination file. Its parent directory must exist. 

2560 overwrite : `bool` 

2561 If True and a file exists at `destination_url` it will be 

2562 overwritten. Otherwise an exception is raised. 

2563 """ 

2564 # Check the source is a file 

2565 if self.is_dir(source_url): 

2566 raise NotImplementedError(f"copy is not implemented for directory {source_url}") 

2567 

2568 # Neither dCache nor XrootD currently implement the COPY 

2569 # webDAV method as documented in 

2570 # 

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

2572 # 

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

2574 # With those servers use third-party copy instead. 

2575 return self._copy_via_third_party(source_url, destination_url, overwrite) 

2576 

2577 def _copy_via_third_party(self, source_url: str, destination_url: str, overwrite: bool = False) -> None: 

2578 """Copy the file at `source_url` to `destination_url` in the same 

2579 storage endpoint using the third-party copy functionality 

2580 implemented by dCache and XRootD servers. 

2581 

2582 Parameters 

2583 ---------- 

2584 source_url : `str` 

2585 URL of the source file. 

2586 destination_url : `str` 

2587 URL of the destination file. Its parent directory must exist. 

2588 overwrite : `bool` 

2589 If True and a file exists at `destination_url` it will be 

2590 overwritten. Otherwise an exception is raised. 

2591 """ 

2592 # To implement COPY we use dCache's third-party copy mechanism 

2593 # documented at: 

2594 # 

2595 # https://www.dcache.org/manuals/UserGuide-10.2/webdav.shtml#third-party-transfers 

2596 # 

2597 # The reason is that dCache does not correctly implement webDAV's COPY 

2598 # method. See https://github.com/dCache/dcache/issues/6950 

2599 

2600 # Create the destination's parent directory first because COPY may 

2601 # fail if it does not exist, depending on the server implementation 

2602 # of RFC 4918. 

2603 destination_parent = self._parent(destination_url) 

2604 self.mkcol(destination_parent) 

2605 

2606 # Retrieve a macaroon for downloading the source 

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

2608 

2609 # Prepare and send the COPY request 

2610 try: 

2611 headers = { 

2612 "Source": source_url, 

2613 "TransferHeaderAuthorization": f"Bearer {download_macaroon}", 

2614 "Credential": "none", 

2615 "Depth": "0", 

2616 "Overwrite": "T" if overwrite else "F", 

2617 "RequireChecksumVerification": "false", 

2618 } 

2619 resp = self._copy(destination_url, headers=headers, preload_content=False) 

2620 match resp.status: 

2621 case HTTPStatus.CREATED: 

2622 return 

2623 case HTTPStatus.ACCEPTED: 

2624 pass 

2625 case _: 

2626 raise ValueError( 

2627 f"Unable to copy resource {resp.geturl()}; status: {resp.status} {resp.reason}" 

2628 ) 

2629 

2630 # Analyse the response to the COPY request that the server has 

2631 # not completed yet. 

2632 content_type = resp.headers.get("Content-Type") 

2633 if content_type != "text/perf-marker-stream": 

2634 raise ValueError( 

2635 f"""Unexpected Content-Type {content_type} in response to COPY request from """ 

2636 f"""{source_url} to {destination_url}""" 

2637 ) 

2638 

2639 # Read the performance markers in the response body until we get 

2640 # a "success" or "failure" notification. 

2641 # 

2642 # Documentation: 

2643 # https://dcache.org/manuals/UserGuide-10.2/webdav.shtml#third-party-transfers 

2644 for marker in io.TextIOWrapper(resp): # type: ignore 

2645 marker = marker.rstrip("\n") 

2646 if marker == "": # EOF 

2647 raise ValueError( 

2648 f"""Copying file from {source_url} to {destination_url} failed: """ 

2649 """could not get response from server""" 

2650 ) 

2651 elif marker.startswith("failure:"): 

2652 raise ValueError( 

2653 f"""Copying file from {source_url} to {destination_url} failed with error: """ 

2654 f"""{marker}""" 

2655 ) 

2656 elif marker.startswith("success:"): 

2657 return 

2658 finally: 

2659 resp.drain_conn() 

2660 

2661 

2662class DavClientDCache(DavClientURLSigner): 

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

2664 

2665 Instances of this class are thread-safe. 

2666 

2667 Parameters 

2668 ---------- 

2669 url : `str` 

2670 Root URL of the storage endpoint 

2671 (e.g. "https://host.example.org:1234/"). 

2672 config : `DavConfig` 

2673 Configuration to initialize this client. 

2674 accepts_ranges : `bool` | `None` 

2675 Indicate whether the remote server accepts the ``Range`` header in GET 

2676 requests. 

2677 """ 

2678 

2679 # Regular expression to parse dCache's response body of a successful 

2680 # PUT request. Such a response body is of the form: 

2681 # 

2682 # "104857600 bytes uploaded\r\n\r\n" 

2683 # 

2684 rex: re.Pattern = re.compile(r"^(\d*) bytes uploaded", re.IGNORECASE | re.ASCII) 

2685 

2686 def __init__(self, url: str, config: DavConfig, accepts_ranges: bool | None = None) -> None: 

2687 super().__init__(url=url, config=config, accepts_ranges=accepts_ranges) 

2688 

2689 # Create a specialized pool manager for sending requests to dCache 

2690 # webdav door, in particular for retrieving metadata. 

2691 # 

2692 # As of dCache v10.2.14, the webDAV door leaves the network connection 

2693 # unusable for us for sending subsequent requests after serving 

2694 # GET, PUT, DELETE, etc., but leaves the connection intact after 

2695 # serving MKCOL, MOVE and PROPFIND requests. 

2696 # We take advantage of that by using a dedicated pool manager for 

2697 # those requests, so that the network connections managed by that pool 

2698 # be reused. This avoids establishing the TCP+TLS connection for each 

2699 # request. 

2700 pool_manager = self._make_pool_manager(self._config) 

2701 self._propfind_pool_manager = pool_manager 

2702 self._move_pool_manager = pool_manager 

2703 self._mkcol_pool_manager = pool_manager 

2704 

2705 # dCache does not deliver macaroons when we are not using a secure 

2706 # channel to interact with the door. In that case, we can not use 

2707 # third party copy and dCache does not correctly support the COPY 

2708 # method as stated in RFC-4918. 

2709 self._can_duplicate = self._base_url.startswith("https://") 

2710 

2711 @override 

2712 def _mkcol( 

2713 self, 

2714 url: str, 

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

2716 pool_manager: PoolManager | None = None, 

2717 ) -> HTTPResponse: 

2718 # Docstring inherited. 

2719 return self._request("MKCOL", url=url, headers=headers, pool_manager=self._mkcol_pool_manager) 

2720 

2721 @override 

2722 def _move( 

2723 self, 

2724 url: str, 

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

2726 pool_manager: PoolManager | None = None, 

2727 ) -> HTTPResponse: 

2728 # Docstring inherited. 

2729 return self._request("MOVE", url=url, headers=headers, pool_manager=self._move_pool_manager) 

2730 

2731 @override 

2732 def _propfind( 

2733 self, 

2734 url: str, 

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

2736 body: str = "", 

2737 pool_manager: PoolManager | None = None, 

2738 ) -> HTTPResponse: 

2739 # Docstring inherited. 

2740 return self._request( 

2741 "PROPFIND", url=url, headers=headers, body=body, pool_manager=self._propfind_pool_manager 

2742 ) 

2743 

2744 @override 

2745 def put( 

2746 self, 

2747 url: str, 

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

2749 data: BinaryIO | bytes = b"", 

2750 ) -> int | None: 

2751 # Docstring inherited. 

2752 

2753 # Send a PUT request with empty body to the dCache frontend server to 

2754 # get redirected to the backend. 

2755 # 

2756 # Details: 

2757 # https://www.dcache.org/manuals/UserGuide-10.2/webdav.shtml#redirection 

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

2759 frontend_headers.update({"Content-Length": "0", "Expect": "100-continue"}) 

2760 if is_zero_length := isinstance(data, bytes) and len(data) == 0: 

2761 # We are uploading an empty file. Don't send the "Expect" header 

2762 # so that the dCache door handles this PUT request itself without 

2763 # redirecting us to a pool. 

2764 frontend_headers.pop("Expect") 

2765 

2766 resp = self._put(url, headers=frontend_headers, body=b"", redirect=False) 

2767 match resp.status: 

2768 case HTTPStatus.OK | HTTPStatus.CREATED | HTTPStatus.NO_CONTENT: 

2769 redirect_url = url 

2770 case status if status in resp.REDIRECT_STATUSES: 

2771 redirect_url = resp.headers.get("Location") 

2772 case _: 

2773 raise unexpected_status_error("PUT", url, resp) 

2774 

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

2776 if is_zero_length: 

2777 return 0 

2778 

2779 # We may have beend redirected to a backend server. Upload the file 

2780 # contents to its final destination. Explicitly ask the server to close 

2781 # this network connection after serving this PUT request to release 

2782 # the associated dCache mover. 

2783 backend_headers = {} if headers is None else dict(headers) 

2784 backend_headers.update({"Connection": "close"}) 

2785 

2786 # Ask dCache to compute and record a checksum of the uploaded 

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

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

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

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

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

2792 # 

2793 # See RFC-3230 for details and 

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

2795 # for the list of supported digest algorithhms. 

2796 if (checksum := self._config.request_checksum) is not None: 

2797 backend_headers.update({"Want-Digest": checksum}) 

2798 

2799 resp = self._put(redirect_url, body=data, headers=backend_headers) 

2800 match resp.status: 

2801 case HTTPStatus.OK | HTTPStatus.CREATED | HTTPStatus.NO_CONTENT: 

2802 # Parse the response body and extract the number of bytes 

2803 # uploaded. This allows us to avoid sending a HEAD request 

2804 # to retrieve the file size. 

2805 response_body = resp.data.decode() 

2806 if match := DavClientDCache.rex.match(response_body): 

2807 return int(match.group(1)) 

2808 else: 

2809 return None 

2810 case _: 

2811 raise unexpected_status_error("PUT", redirect_url, resp) 

2812 

2813 @override 

2814 def download(self, url: str, filename: str, chunk_size: int) -> int: 

2815 """Download the content of a file and write it to local file. 

2816 

2817 Parameters 

2818 ---------- 

2819 url : `str` 

2820 Target URL. 

2821 filename : `str` 

2822 Local file to write the content to. If the file already exists, 

2823 it will be rewritten. 

2824 chunk_size : `int` 

2825 Size of the chunks to write to `filename`. 

2826 

2827 Returns 

2828 ------- 

2829 count: `int` 

2830 Number of bytes written to `filename`. 

2831 

2832 Notes 

2833 ----- 

2834 The caller must ensure that the resource at `url` is a file, not 

2835 a directory. 

2836 """ 

2837 # Send a GET request without following redirection to get redirected 

2838 # to the backend server. 

2839 _, resp = self.get(url, preload_content=False, redirect=False) 

2840 match resp.status: 

2841 case HTTPStatus.OK: 

2842 # We were not redirected. Consume this response. 

2843 return self._write_response_body_to_file(resp, filename, chunk_size) 

2844 case status if status not in resp.REDIRECT_STATUSES: 

2845 raise unexpected_status_error("GET", url, resp) 

2846 case _: 

2847 # We were redirected. Follow this redirection. 

2848 pass 

2849 

2850 # Drain and release the response we received from the frontend server 

2851 # so that the connection can be reused. 

2852 resp.drain_conn() 

2853 resp.release_conn() 

2854 

2855 # We were redirected to a backend server. Send a GET request to the 

2856 # backend server and ask it to close the HTTP connection to force 

2857 # closing the network connection. 

2858 redirect_url = resp.headers.get("Location") 

2859 _, resp = self.get(redirect_url, headers={"Connection": "close"}, preload_content=False) 

2860 match resp.status: 

2861 case HTTPStatus.OK: 

2862 return self._write_response_body_to_file(resp, filename, chunk_size) 

2863 case _: 

2864 raise unexpected_status_error("GET", redirect_url, resp) 

2865 

2866 @override 

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

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

2869 

2870 Parameters 

2871 ---------- 

2872 url : `str` 

2873 Target URL. 

2874 

2875 Returns 

2876 ------- 

2877 url: `str` 

2878 Backend URL from which the data was obtained. 

2879 data: `bytes` 

2880 Contents of the file. 

2881 

2882 Notes 

2883 ----- 

2884 The caller must ensure that the resource at `url` is a file, not 

2885 a directory. 

2886 """ 

2887 # Send a GET request without following redirection to get redirected 

2888 # to the backend server. 

2889 backend_url, resp = self.get(url, redirect=False) 

2890 match resp.status: 

2891 case HTTPStatus.OK: 

2892 return backend_url, resp.data 

2893 case status if status in resp.REDIRECT_STATUSES: 

2894 redirect_url = resp.headers.get("Location") 

2895 case _: 

2896 raise unexpected_status_error("GET", url, resp) 

2897 

2898 # We were redirected. Send a GET request to the backend server 

2899 # and ask it to close the HTTP connection to force closing the 

2900 # network connection. 

2901 final_url, resp = self.get(redirect_url, headers={"Connection": "close"}) 

2902 match resp.status: 

2903 case HTTPStatus.OK: 

2904 return final_url, resp.data 

2905 case _: 

2906 raise unexpected_status_error("GET", redirect_url, resp) 

2907 

2908 @override 

2909 def write(self, url: str, data: BinaryIO | bytes) -> int | None: 

2910 """Create or rewrite a remote file at `url` with `data` as its 

2911 contents. 

2912 

2913 Parameters 

2914 ---------- 

2915 url : `str` 

2916 Target URL. 

2917 data : `bytes` 

2918 Sequence of bytes to upload. 

2919 

2920 Returns 

2921 ------- 

2922 size : `int | None` 

2923 The size in bytes of the file uploaded. Can be `None` if the size 

2924 could not be retrieved. 

2925 

2926 Notes 

2927 ----- 

2928 If a file already exists at `url` it will be rewritten. 

2929 """ 

2930 # dCache will automatically create all the parent directories so we 

2931 # don't need to explicitly create them. Although this is not compliant 

2932 # to RFC 4918, this is advantageous because it avoids several 

2933 # round-trips to the server for creating all the directories 

2934 # before actually uploading the data. 

2935 try: 

2936 # Upload to a temporary file and rename to the final name. 

2937 temporary_url = self._make_temporary_url(url) 

2938 size = self.put(temporary_url, data=data) 

2939 self.rename(temporary_url, url, overwrite=True, create_parent=False) 

2940 

2941 # Update the file size cache with this size 

2942 self._file_size_cache.update_size(url, size) 

2943 return size 

2944 except Exception: 

2945 # Upload failed. Attempt to remove the temporary file. 

2946 self.delete(temporary_url) 

2947 raise 

2948 

2949 @override 

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

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

2952 

2953 If a directory already exists at `url` no error is returned nor 

2954 exception is raised. An exception is raised if a file exists at `url`. 

2955 

2956 Parameters 

2957 ---------- 

2958 url : `str` 

2959 Target URL. 

2960 """ 

2961 # A "MKCOL" request to dCache does not automatically create all 

2962 # the intermediate directories if they do not exist. However, a 

2963 # "PUT" request of a file does create the directory hierarchy. 

2964 # 

2965 # We exploit that to create directory hierarchies: we first create an 

2966 # empty file with a random name and then we remove it. As a side 

2967 # effect, the target directory will be created. 

2968 # 

2969 # Creating a directory this way implies two requests to the server 

2970 # ("PUT" and "DELETE"), while using "MKCOL" would on average imply 

2971 # one request per inexisting directory in the hierarchy. When 

2972 # directory hierarchies are relatively deep, requiring two 

2973 # requests per hierarchy is better than sending a "MKCOL" request 

2974 # per directory in the hierarchy. 

2975 try: 

2976 temporary_url = self._make_temporary_url(url=f"{url}/mkcol") 

2977 self.put(temporary_url, data=b"") 

2978 finally: 

2979 self.delete(temporary_url) 

2980 

2981 @override 

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

2983 # Docstring inherited. 

2984 result: dict[str, Any] = { 

2985 "name": name if name is not None else url, 

2986 "type": None, 

2987 "size": None, 

2988 "last_modified": datetime.min, 

2989 "checksums": {}, 

2990 } 

2991 

2992 # Request live DAV properties as well as the checksums that dCache 

2993 # recorded about this file. 

2994 body = ( 

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

2996 """<D:propfind xmlns:D="DAV:" xmlns:dcache="http://www.dcache.org/2013/webdav">""" 

2997 """<D:prop>""" 

2998 """<D:resourcetype/>""" 

2999 """<D:getcontentlength/>""" 

3000 """<D:getlastmodified/>""" 

3001 """<D:displayname/>""" 

3002 """<dcache:Checksums/>""" 

3003 """</D:prop>""" 

3004 """</D:propfind>""" 

3005 ) 

3006 resp = self.propfind(url, body=body, depth="0") 

3007 match resp.status: 

3008 case HTTPStatus.NOT_FOUND: 

3009 return result 

3010 case HTTPStatus.MULTI_STATUS: 

3011 property = self._propfind_parser.parse(resp.data)[0] 

3012 metadata = DavFileMetadata.from_property(base_url=self._base_url, property=property) 

3013 result.update( 

3014 { 

3015 "type": "directory" if metadata.is_dir else "file", 

3016 "size": metadata.size, 

3017 "last_modified": metadata.last_modified, 

3018 "checksums": metadata.checksums, 

3019 } 

3020 ) 

3021 return result 

3022 case _: 

3023 raise unexpected_status_error("PROPFIND", url, resp) 

3024 

3025 @override 

3026 def read_range( 

3027 self, 

3028 url: str, 

3029 start: int, 

3030 end: int | None, 

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

3032 ) -> tuple[str, bytes]: 

3033 # Docstring inherited. 

3034 range_headers = {"Accept-Encoding": "identity"} 

3035 if end is None: 

3036 range_headers.update({"Range": f"bytes={start}-"}) 

3037 else: 

3038 range_headers.update({"Range": f"bytes={start}-{end}"}) 

3039 

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

3041 frontend_headers.update(range_headers) 

3042 

3043 # Send the GET request to the dCache door but don't follow redirections 

3044 # automatically. We need to be able to add a `Connection: close` 

3045 # request header when sending the request to the dCache pool we 

3046 # will be redirected to. 

3047 # 

3048 # We don't send that header to the door since we want to keep the 

3049 # network connection with the door open for later reuse. 

3050 final_url, resp = self.get(url, headers=frontend_headers, redirect=False) 

3051 match resp.status: 

3052 case HTTPStatus.PARTIAL_CONTENT: 

3053 return final_url, resp.data 

3054 case status if status not in resp.REDIRECT_STATUSES: 

3055 raise unexpected_status_error("GET (with 'Range' header)", url, resp) 

3056 case _: 

3057 pass 

3058 

3059 # We were redirected to the dCache pool. Follow the redirection and 

3060 # add a `Connection: close` header to notify the mover to stop its 

3061 # execution after serving this request. 

3062 backend_headers = {} if headers is None else dict(headers) 

3063 backend_headers.update(range_headers) 

3064 backend_headers.update({"Connection": "close"}) 

3065 

3066 redirect_url = resp.headers.get("Location") 

3067 _, resp = self.get(redirect_url, headers=backend_headers, redirect=True) 

3068 match resp.status: 

3069 case HTTPStatus.PARTIAL_CONTENT: 

3070 # Return the door URL so that subsequent requests (if any) 

3071 # go through the dCache door instead first of going directly to 

3072 # the pool since we asked the mover to be stopped. 

3073 return url, resp.data 

3074 case _: 

3075 raise unexpected_status_error("GET (with 'Range' header)", redirect_url, resp) 

3076 

3077 @override 

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

3079 # Docstring inherited. 

3080 

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

3082 # connection with the dCache pool open. 

3083 pass 

3084 

3085 

3086class DavClientXrootD(DavClientURLSigner): 

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

3088 

3089 Instances of this class are thread-safe. 

3090 

3091 Parameters 

3092 ---------- 

3093 url : `str` 

3094 Root URL of the storage endpoint 

3095 (e.g. "https://host.example.org:1234/"). 

3096 config : `DavConfig` 

3097 Configuration to initialize this client. 

3098 accepts_ranges : `bool` | `None` 

3099 Indicate whether the remote server accepts the ``Range`` header in GET 

3100 requests. 

3101 """ 

3102 

3103 def __init__(self, url: str, config: DavConfig, accepts_ranges: bool | None = None) -> None: 

3104 super().__init__(url=url, config=config, accepts_ranges=accepts_ranges) 

3105 

3106 @override 

3107 def put( 

3108 self, 

3109 url: str, 

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

3111 data: BinaryIO | bytes = b"", 

3112 ) -> int | None: 

3113 # Docstring inherited. 

3114 

3115 # Send a PUT request with empty body to the XRootD frontend server to 

3116 # get redirected to the backend. 

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

3118 frontend_headers.update({"Content-Length": "0", "Expect": "100-continue"}) 

3119 for attempt in range(max_attempts := 3): 

3120 resp = self._put(url, headers=frontend_headers, body=b"", redirect=False) 

3121 if resp.status in ( 

3122 HTTPStatus.OK, 

3123 HTTPStatus.CREATED, 

3124 HTTPStatus.NO_CONTENT, 

3125 ): 

3126 redirect_url = url 

3127 break 

3128 elif resp.status in resp.REDIRECT_STATUSES: 

3129 redirect_url = resp.headers.get("Location") 

3130 break 

3131 elif resp.status == HTTPStatus.LOCKED: 

3132 # Sometimes XRootD servers respond with status code LOCKED and 

3133 # response body of the form: 

3134 # 

3135 # "Output file /path/to/file is already opened by 1 writer; 

3136 # open denied." 

3137 # 

3138 # If we get such a response, try again, unless we reached 

3139 # the maximum number of attempts. 

3140 if attempt == max_attempts - 1: 

3141 raise ValueError( 

3142 f"""Unexpected response to HTTP request PUT {resp.geturl()}: status {resp.status} """ 

3143 f"""{resp.reason} [{resp.data.decode()}] after {max_attempts} attempts""" 

3144 ) 

3145 

3146 # Wait a bit and try again 

3147 log.warning( 

3148 f"""got unexpected response status {HTTPStatus.LOCKED} Locked for PUT {resp.geturl()} """ 

3149 f"""(attempt {attempt}/{max_attempts}), retrying...""" 

3150 ) 

3151 time.sleep((attempt + 1) * 0.100) 

3152 continue 

3153 else: 

3154 raise unexpected_status_error("PUT", url, resp) 

3155 

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

3157 # its final destination. 

3158 

3159 # XRootD backend servers typically use a single port number for 

3160 # accepting connections from clients. It is therefore beneficial 

3161 # to keep those connections open, if the server allows. 

3162 

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

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

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

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

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

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

3169 # 

3170 # See RFC-3230 for details and 

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

3172 # for the list of supported digest algorithhms. 

3173 # 

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

3175 # the checksum reqquest may be ignored by the server. 

3176 backend_headers = {} if headers is None else dict(headers) 

3177 if (checksum := self._config.request_checksum) is not None: 

3178 backend_headers.update({"Want-Digest": checksum}) 

3179 

3180 resp = self._put(redirect_url, body=data, headers=backend_headers) 

3181 match resp.status: 

3182 case HTTPStatus.OK | HTTPStatus.CREATED | HTTPStatus.NO_CONTENT: 

3183 # Send a HEAD request to retrieve the size of the file we 

3184 # just uploaded. 

3185 resp = self.head(redirect_url) 

3186 size = int(resp.headers.get("Content-Length", -1)) 

3187 return None if size == -1 else size 

3188 case _: 

3189 raise unexpected_status_error("PUT", redirect_url, resp) 

3190 

3191 @override 

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

3193 # XRootD does not include checksums in the response to PROPFIND 

3194 # request. We need to send a specific HEAD request to retrieve 

3195 # the ADLER32 checksum. 

3196 # 

3197 # If found, the checksum is included in the response header "Digest", 

3198 # which is of the form: 

3199 # 

3200 # Digest: adler32=0e4709f2 

3201 result = super().info(url, name) 

3202 if result["type"] == "file": 

3203 headers: dict[str, str] = {"Want-Digest": "adler32"} 

3204 resp = self.head(url=url, headers=headers) 

3205 if (digest := resp.headers.get("Digest")) is not None: 

3206 value = digest.split("=")[1] 

3207 result["checksums"].update({"adler32": value}) 

3208 

3209 return result 

3210 

3211 @override 

3212 def write(self, url: str, data: BinaryIO | bytes) -> int | None: 

3213 """Create or rewrite a remote file at `url` with `data` as its 

3214 contents. 

3215 

3216 Parameters 

3217 ---------- 

3218 url : `str` 

3219 Target URL. 

3220 data : `bytes` 

3221 Sequence of bytes to upload. 

3222 

3223 Returns 

3224 ------- 

3225 size : `int | None` 

3226 The size in bytes of the file uploaded. Can be `None` if the size 

3227 could not be retrieved. 

3228 

3229 Notes 

3230 ----- 

3231 If a file already exists at `url` it will be rewritten. 

3232 """ 

3233 # XRootD will automatically create all the parent directories so we 

3234 # don't need to explicitly create them. Although this is not compliant 

3235 # to RFC 4918, this is advantageous because it avoids several 

3236 # round-trips to the server for creating all the directories 

3237 # before actually uploading the data. 

3238 try: 

3239 # Upload to a temporary file and rename to the final name. 

3240 temporary_url = self._make_temporary_url(url) 

3241 size = self.put(temporary_url, data=data) 

3242 self.rename(temporary_url, url, overwrite=True, create_parent=False) 

3243 

3244 # Update the file size cache with this size 

3245 self._file_size_cache.update_size(url, size) 

3246 return size 

3247 except Exception: 

3248 # Upload failed. Attempt to remove the temporary file. 

3249 self.delete(temporary_url) 

3250 raise 

3251 

3252 @override 

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

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

3255 

3256 If a directory already exists at `url` no error is returned nor 

3257 exception is raised. An exception is raised if a file exists at `url`. 

3258 

3259 Parameters 

3260 ---------- 

3261 url : `str` 

3262 Target URL. 

3263 """ 

3264 # XRootD automatically creates all the intermediate directories. 

3265 resp = self._mkcol(url) 

3266 match resp.status: 

3267 case HTTPStatus.CREATED: 

3268 return 

3269 case HTTPStatus.METHOD_NOT_ALLOWED: 

3270 # XRootD returns "405 Method Not Allowed" when either a file 

3271 # or a directory already exists at `url` 

3272 stat = self.stat(url) 

3273 if stat.is_dir: 

3274 # A directory exists at `url`. Nothing more to do. 

3275 return 

3276 elif stat.is_file: 

3277 raise NotADirectoryError( 

3278 f"Can not create a directory because a file already exists at {resp.geturl()}" 

3279 ) 

3280 case _: 

3281 raise ValueError( 

3282 f"Can not create directory {resp.geturl()}: status {resp.status} {resp.reason}" 

3283 ) 

3284 

3285 @override 

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

3287 # Docstring inherited. 

3288 

3289 # XRootD v5.9.1 responds "200 OK" to a HEAD request against an 

3290 # existing file. When the target URL is a directory, it also responds 

3291 # "200 OK". In both cases the response header "Content-Length" 

3292 # is present but has different meaning. If the target URL is a file, 

3293 # the header value is the size in bytes of the file. If the target 

3294 # URL is a directory, the header value is the number of items in 

3295 # the directory. 

3296 # 

3297 # So there is not an easy way to determine if the target URL is a 

3298 # file or a directory from the response to a HEAD request. 

3299 # 

3300 # When the target URL is a directory and we ask for a digest, the 

3301 # server responds "409 Conflict". We use this behavior to 

3302 # discriminate between a file and a directory. 

3303 # 

3304 # Note that XRootD does not include the "Last-Modified" header in the 

3305 # response to a HEAD request so we cannot include the last modified 

3306 # time in the value returned by this method. 

3307 resp = self._head(url, headers={"Want-Digest": "adler32"}) 

3308 match resp.status: 

3309 case HTTPStatus.OK: 

3310 # There is a file at target URL 

3311 if "Content-Length" in resp.headers: 

3312 href = url.replace(self._base_url, "", 1) 

3313 size = int(resp.headers.get("Content-Length")) 

3314 return DavFileMetadata(self._base_url, href=href, exists=True, is_dir=False, size=size) 

3315 else: 

3316 raise ValueError( 

3317 f"""Expecting Content-Length header to be present in """ 

3318 f"""response to HTTP HEAD {resp.geturl()}: status {resp.status} """ 

3319 f"""{resp.reason} [{resp.data.decode()}] but could not find it""" 

3320 ) 

3321 case HTTPStatus.CONFLICT: 

3322 # There is a directory at target URL 

3323 href = url.replace(self._base_url, "", 1) 

3324 return DavFileMetadata(self._base_url, href=href, exists=True, is_dir=True) 

3325 case HTTPStatus.NOT_FOUND: 

3326 # There is neither a file nor a directory at target URL 

3327 return DavFileMetadata(base_url=url, exists=False) 

3328 case _: 

3329 raise unexpected_status_error("HEAD", url, resp) 

3330 

3331 @override 

3332 def read_range( 

3333 self, 

3334 url: str, 

3335 start: int, 

3336 end: int | None, 

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

3338 ) -> tuple[str, bytes]: 

3339 # Docstring inherited. 

3340 

3341 # Send the request to the XRootD redirector and follow 

3342 # redirections automatically. 

3343 # 

3344 # The network connection with the redirector and with the backend 

3345 # file server are left open for later reuse. 

3346 range_headers = {"Accept-Encoding": "identity"} 

3347 if end is None: 

3348 range_headers.update({"Range": f"bytes={start}-"}) 

3349 else: 

3350 range_headers.update({"Range": f"bytes={start}-{end}"}) 

3351 

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

3353 get_headers.update(range_headers) 

3354 

3355 final_url, resp = self.get(url, headers=get_headers, redirect=True) 

3356 match resp.status: 

3357 case HTTPStatus.PARTIAL_CONTENT: 

3358 return final_url, resp.data 

3359 case _: 

3360 raise unexpected_status_error("GET (with 'Range' header)", url, resp) 

3361 

3362 @override 

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

3364 # Docstring inherited. 

3365 

3366 # Send a `HEAD` request with a `Connection: close` header to notify 

3367 # the remote server that we are not sending other GET requests with 

3368 # `Range` header for this URL. 

3369 self._request("HEAD", url=url, headers={"Connection": "close"}, redirect=False) 

3370 

3371 

3372class DavFileMetadata: 

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

3374 

3375 Parameters 

3376 ---------- 

3377 base_url : `str` 

3378 Base URL. 

3379 href : `str`, optional 

3380 Path component that can be added to the base URL. 

3381 name : `str`, optional 

3382 Name. 

3383 exists : `bool`, optional 

3384 Whether file or directory exist. 

3385 size : `int`, optional 

3386 Size of file. 

3387 is_dir : `bool`, optional 

3388 Whether the URL points to a directory or file. 

3389 last_modified : `bool`, optional 

3390 Last modified date. 

3391 checksums : `dict` [ `str`, `str` ] | `None`, optional 

3392 Checksums. 

3393 """ 

3394 

3395 def __init__( 

3396 self, 

3397 base_url: str, 

3398 href: str = "", 

3399 name: str = "", 

3400 exists: bool = False, 

3401 size: int = -1, 

3402 is_dir: bool = False, 

3403 last_modified: datetime = datetime.min, 

3404 checksums: dict[str, str] | None = None, 

3405 ): 

3406 self._url: str = base_url if not href else base_url.rstrip("/") + href 

3407 self._href: str = href 

3408 self._name: str = name 

3409 self._exists: bool = exists 

3410 self._size: int = size 

3411 self._is_dir: bool = is_dir 

3412 self._last_modified: datetime = last_modified 

3413 self._checksums: dict[str, str] = {} if checksums is None else dict(checksums) 

3414 

3415 @staticmethod 

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

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

3418 

3419 Parameters 

3420 ---------- 

3421 base_url : `str` 

3422 Base URL. 

3423 property : `DavProperty` 

3424 Properties to associate with URL. 

3425 """ 

3426 return DavFileMetadata( 

3427 base_url=base_url, 

3428 href=property.href, 

3429 name=property.name, 

3430 exists=property.exists, 

3431 size=property.size, 

3432 is_dir=property.is_dir, 

3433 last_modified=property.last_modified, 

3434 checksums=dict(property.checksums), 

3435 ) 

3436 

3437 def __str__(self) -> str: 

3438 return ( 

3439 f"""{self._url} {self._href} {self._name} {self._exists} {self._size} {self._is_dir} """ 

3440 f"""{self._checksums}""" 

3441 ) 

3442 

3443 @property 

3444 def url(self) -> str: 

3445 return self._url 

3446 

3447 @property 

3448 def href(self) -> str: 

3449 return self._href 

3450 

3451 @property 

3452 def name(self) -> str: 

3453 return self._name 

3454 

3455 @property 

3456 def exists(self) -> bool: 

3457 return self._exists 

3458 

3459 @property 

3460 def size(self) -> int: 

3461 if not self._exists: 

3462 return -1 

3463 

3464 return 0 if self._is_dir else self._size 

3465 

3466 @property 

3467 def is_dir(self) -> bool: 

3468 return self._exists and self._is_dir 

3469 

3470 @property 

3471 def is_file(self) -> bool: 

3472 return self._exists and not self._is_dir 

3473 

3474 @property 

3475 def last_modified(self) -> datetime: 

3476 return self._last_modified 

3477 

3478 @property 

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

3480 return self._checksums 

3481 

3482 

3483class DavProperty: 

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

3485 resource, as retrieved via a PROPFIND request. 

3486 

3487 Parameters 

3488 ---------- 

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

3490 The XML response defining the DAV property. 

3491 """ 

3492 

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

3494 # PROPFIND response's 'propstat' element. 

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

3496 

3497 def __init__(self, response: eTree.Element | None): 

3498 self._href: str = "" 

3499 self._displayname: str = "" 

3500 self._collection: bool = False 

3501 self._getlastmodified: str = "" 

3502 self._getcontentlength: int = -1 

3503 self._checksums: dict[str, str] = {} 

3504 

3505 if response is not None: 

3506 self._parse(response) 

3507 

3508 def _parse(self, response: eTree.Element) -> None: 

3509 # Extract 'href'. 

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

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

3512 # keep mypy happy. 

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

3514 else: 

3515 raise ValueError( 

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

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

3518 ) 

3519 

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

3521 # Only extract properties of interest with status OK. 

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

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

3524 continue 

3525 

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

3527 # Parse "collection". 

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

3529 self._collection = True 

3530 

3531 # Parse "getlastmodified". 

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

3533 self._getlastmodified = str(element.text) 

3534 

3535 # Parse "getcontentlength". 

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

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

3538 

3539 # Parse "displayname". 

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

3541 self._displayname = str(element.text) 

3542 

3543 # Parse "Checksums" 

3544 if (element := prop.find("./{http://www.dcache.org/2013/webdav}Checksums")) is not None: 

3545 self._checksums = self._parse_checksums(element.text) 

3546 

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

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

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

3550 if not self._displayname: 

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

3552 

3553 # Some webDAV servers do not append a "/" to the href of directories. 

3554 # Ensure we include a single final "/" in our response. 

3555 if self._collection: 

3556 self._href = self._href.rstrip("/") + "/" 

3557 

3558 # Force a size of 0 for collections. 

3559 if self._collection: 

3560 self._getcontentlength = 0 

3561 

3562 def _parse_checksums(self, checksums: str | None) -> dict[str, str]: 

3563 # checksums argument is of the form 

3564 # md5=MyS/wljSzI9WYiyrsuyoxw==,adler32=23b104f2 

3565 result: dict[str, str] = {} 

3566 if checksums is not None: 

3567 for checksum in checksums.split(","): 

3568 if (pos := checksum.find("=")) != -1: 

3569 algorithm, value = (checksum[:pos].lower(), checksum[pos + 1 :]) 

3570 if algorithm == "md5": 

3571 # dCache documentation about how it encodes the 

3572 # MD5 checksum: 

3573 # 

3574 # https://www.dcache.org/manuals/UserGuide-10.2/webdav.shtml#checksums 

3575 result[algorithm] = bytes.hex(base64.standard_b64decode(value)) 

3576 else: 

3577 result[algorithm] = value 

3578 

3579 return result 

3580 

3581 @property 

3582 def exists(self) -> bool: 

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

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

3585 

3586 @property 

3587 def is_dir(self) -> bool: 

3588 return self._collection 

3589 

3590 @property 

3591 def is_file(self) -> bool: 

3592 return not self._collection 

3593 

3594 @property 

3595 def last_modified(self) -> datetime: 

3596 if not self._getlastmodified: 

3597 return datetime.min 

3598 

3599 # Last modified timestamp is of the form: 

3600 # 'Wed, 12 Mar 2025 10:11:13 GMT' 

3601 return datetime.strptime(self._getlastmodified, "%a, %d %b %Y %H:%M:%S %Z").replace(tzinfo=UTC) 

3602 

3603 @property 

3604 def size(self) -> int: 

3605 return self._getcontentlength 

3606 

3607 @property 

3608 def name(self) -> str: 

3609 return self._displayname 

3610 

3611 @property 

3612 def href(self) -> str: 

3613 return self._href 

3614 

3615 @property 

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

3617 return self._checksums 

3618 

3619 

3620class DavPropfindParser: 

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

3622 

3623 def __init__(self) -> None: 

3624 return 

3625 

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

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

3628 PROPFIND request. 

3629 

3630 Parameters 

3631 ---------- 

3632 body : `bytes` 

3633 XML-encoded response body to a PROPFIND request. 

3634 

3635 Returns 

3636 ------- 

3637 responses : `list` [ `DavProperty` ] 

3638 Parsed content of the response. 

3639 

3640 Notes 

3641 ----- 

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

3643 this function raises. 

3644 """ 

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

3646 # readability): 

3647 # 

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

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

3650 # <D:response> 

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

3652 # <D:propstat> 

3653 # <D:prop> 

3654 # <D:resourcetype> 

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

3656 # </D:resourcetype> 

3657 # <D:getlastmodified> 

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

3659 # </D:getlastmodified> 

3660 # <D:getcontentlength> 

3661 # 12345 

3662 # </D:getcontentlength> 

3663 # </D:prop> 

3664 # <D:status> 

3665 # HTTP/1.1 200 OK 

3666 # </D:status> 

3667 # </D:propstat> 

3668 # </D:response> 

3669 # <D:response> 

3670 # ... 

3671 # </D:response> 

3672 # <D:response> 

3673 # ... 

3674 # </D:response> 

3675 # </D:multistatus> 

3676 

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

3678 decoded_body: str = body.decode().strip() 

3679 responses = [] 

3680 multistatus = eTree.fromstring(decoded_body) 

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

3682 responses.append(DavProperty(response)) 

3683 

3684 if responses: 

3685 return responses 

3686 else: 

3687 # Could not parse the body 

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

3689 

3690 

3691class Authorizer: 

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

3693 

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

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

3696 

3697 Parameters 

3698 ---------- 

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

3700 Dict to augment with authorization information. 

3701 

3702 Notes 

3703 ----- 

3704 This method must be implemented by concrete subclasses. 

3705 """ 

3706 raise NotImplementedError 

3707 

3708 def _is_file_protected(self, filepath: str) -> bool: 

3709 """Return true if the permissions of file at `filepath` only allow for 

3710 access by its owner. 

3711 

3712 Parameters 

3713 ---------- 

3714 filepath : `str` 

3715 Path of a local file. 

3716 """ 

3717 if not os.path.isfile(filepath): 3717 ↛ 3718line 3717 didn't jump to line 3718 because the condition on line 3717 was never true

3718 return False 

3719 

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

3721 owner_accessible = bool(mode & stat.S_IRWXU) 

3722 group_accessible = bool(mode & stat.S_IRWXG) 

3723 other_accessible = bool(mode & stat.S_IRWXO) 

3724 return owner_accessible and not group_accessible and not other_accessible 

3725 

3726 def _read_if_modified_since( 

3727 self, filename: str | None, timestamp: float 

3728 ) -> tuple[str, float] | tuple[None, None]: 

3729 """Read local file `filename` if its modification time is more 

3730 recent than `timestamp`. 

3731 

3732 Parameters 

3733 ---------- 

3734 filename : `str`, optional 

3735 Path of a local file. 

3736 

3737 timestamp: `float`, optional 

3738 Timestamp to compare against the last modification time of 

3739 `filename`. The contents of file at `filename` is only read if its 

3740 modification time is more recent than `timestamp`. 

3741 

3742 Returns 

3743 ------- 

3744 result: `tuple[str, float]` 

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

3746 operation). 

3747 

3748 If `filename` is `None`, the returned value is `tuple[None, None]`. 

3749 """ 

3750 if filename is None: 3750 ↛ 3751line 3750 didn't jump to line 3751 because the condition on line 3750 was never true

3751 return (None, None) 

3752 

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

3754 return (None, None) 

3755 

3756 with open(filename) as file: 

3757 time_of_last_read = time.time() 

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

3759 

3760 

3761class TokenAuthorizer(Authorizer): 

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

3763 

3764 Parameters 

3765 ---------- 

3766 token : `str` 

3767 Can be either the path to a local file which contains the 

3768 value of the token or the token itself. If `token` is a file 

3769 it must be protected so that only the owner can read and write it. 

3770 """ 

3771 

3772 def __init__(self, token: str | None = None) -> None: 

3773 self._token = self._token_path = None 

3774 self._time_of_last_read: float = -1.0 

3775 if token is None: 

3776 return 

3777 

3778 self._token = token 

3779 if os.path.isfile(token): 

3780 self._token_path = os.path.abspath(token) 

3781 if not self._is_file_protected(self._token_path): 

3782 raise PermissionError( 

3783 f"""Authorization token file at {self._token_path} must be protected for access only """ 

3784 """by its owner""" 

3785 ) 

3786 self._update_token() 

3787 

3788 def _update_token(self) -> None: 

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

3790 than the last time we read it. 

3791 """ 

3792 if self._token_path is None: 

3793 return None 

3794 

3795 token, time_of_last_read = self._read_if_modified_since(self._token_path, self._time_of_last_read) 

3796 if token is None or time_of_last_read is None: 

3797 return 

3798 

3799 # Update the token value and the last time we read it. 

3800 self._token = token 

3801 self._time_of_last_read = time_of_last_read 

3802 

3803 @override 

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

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

3806 

3807 Parameters 

3808 ---------- 

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

3810 Dict to augment with authorization information. 

3811 """ 

3812 if self._token is None: 

3813 return 

3814 

3815 self._update_token() 

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

3817 

3818 

3819class BasicAuthorizer(Authorizer): 

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

3821 authentication. 

3822 

3823 Parameters 

3824 ---------- 

3825 user_name : `str` 

3826 Can be either the path to a local file which contains the 

3827 user name or the user name itself. If `user_name` is a file 

3828 it must be protected so that only the owner can read and write it. 

3829 user_password : `str` 

3830 Can be either the path to a local file which contains the 

3831 value of the password or the password itself. If `user_password` is a 

3832 file it must be protected so that only the owner can read and write it. 

3833 """ 

3834 

3835 def __init__(self, user_name: str | None = None, user_password: str | None = None) -> None: 

3836 if user_name is None or user_password is None: 

3837 return 

3838 

3839 self._user_name: str | None = user_name 

3840 self._user_password: str | None = user_password 

3841 self._user_password_path: str | None = None 

3842 self._time_of_last_read: float = -1.0 

3843 self._header_value: str = "" 

3844 

3845 if os.path.isfile(self._user_password): 

3846 # The value in `user_password` is the path to a file. Check 

3847 # the file is protected and read its contents. 

3848 self._user_password_path = os.path.abspath(self._user_password) 

3849 if not self._is_file_protected(self._user_password_path): 

3850 raise PermissionError( 

3851 f"""Password file at {self._user_password_path} must be protected for access only """ 

3852 """by its owner""" 

3853 ) 

3854 self._update_password() 

3855 else: 

3856 self._update_header_value() 

3857 

3858 def _update_header_value(self) -> None: 

3859 """Compute the value of the 'Authorization' header using HTTP basic 

3860 authorization. 

3861 """ 

3862 basic_auth_header = make_headers(basic_auth=f"{self._user_name}:{self._user_password}") 

3863 self._header_value = basic_auth_header["authorization"] 

3864 

3865 def _update_password(self) -> None: 

3866 """Update the password of this authorizer if the file it is stored in 

3867 has been modified since the last time we read it. 

3868 """ 

3869 if self._user_password_path is None: 

3870 return None 

3871 

3872 password, time_of_last_read = self._read_if_modified_since( 

3873 self._user_password_path, self._time_of_last_read 

3874 ) 

3875 if password is None or time_of_last_read is None: 

3876 return 

3877 

3878 # Update the password, the last time we read it and re-compute the 

3879 # value of the "Authorization" header. 

3880 self._user_password = password 

3881 self._time_of_last_read = time_of_last_read 

3882 self._update_header_value() 

3883 

3884 @override 

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

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

3887 

3888 Parameters 

3889 ---------- 

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

3891 Dict to augment with authorization information. 

3892 """ 

3893 if self._user_name is None or self._user_password is None: 

3894 return 

3895 

3896 self._update_password() 

3897 headers["Authorization"] = self._header_value 

3898 

3899 

3900def expand_vars(path: str | None) -> str | None: 

3901 """Expand the environment variables in `path` and return the path with 

3902 the value of the variable expanded. 

3903 

3904 Parameters 

3905 ---------- 

3906 path : `str` or `None` 

3907 Abolute or relative path which may include an environment variable 

3908 (e.g. '$HOME/path/to/my/file'). 

3909 

3910 Returns 

3911 ------- 

3912 path: `str` 

3913 The path with the values of the environment variables expanded. 

3914 """ 

3915 return None if path is None else os.path.expandvars(path) 

3916 

3917 

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

3919 """Dump response for debugging purposes. 

3920 

3921 Parameters 

3922 ---------- 

3923 method : `str` 

3924 Method name to include in log output. 

3925 resp : `HTTPResponse` 

3926 Response to dump. 

3927 dump_body : `bool`, optional 

3928 Whether or not to issue a debug log message. 

3929 """ 

3930 log.debug("%s %s", method, resp.geturl()) 

3931 log.debug(" %s %s", resp.status, resp.reason) 

3932 

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

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

3935 

3936 if dump_body: 

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