XRootD
Loading...
Searching...
No Matches
XrdClHttpOps.cc
Go to the documentation of this file.
1/******************************************************************************/
2/* Copyright (C) 2025, Pelican Project, Morgridge Institute for Research */
3/* */
4/* This file is part of the XrdClHttp client plugin for XRootD. */
5/* */
6/* XRootD is free software: you can redistribute it and/or modify it under */
7/* the terms of the GNU Lesser General Public License as published by the */
8/* Free Software Foundation, either version 3 of the License, or (at your */
9/* option) any later version. */
10/* */
11/* XRootD is distributed in the hope that it will be useful, but WITHOUT */
12/* ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or */
13/* FITNESS FOR A PARTICULAR PURPOSE. See the GNU Lesser General Public */
14/* License for more details. */
15/* */
16/* The copyright holder's institutional names and contributor's names may not */
17/* be used to endorse or promote products derived from this software without */
18/* specific prior written permission of the institution or contributor. */
19/******************************************************************************/
20
21#include "XrdClHttpOps.hh"
23#include "XrdClHttpResponses.hh"
24#include "XrdClHttpUtil.hh"
25#include "XrdClHttpWorker.hh"
26
28#include <XrdCl/XrdClLog.hh>
30
31#include <arpa/inet.h>
32#include <unistd.h>
33#include <cerrno>
34#include <chrono>
35#include <cmath>
36#ifdef __APPLE__
37#include <stdlib.h>
38#else
39#include <sys/random.h>
40#endif
41#include <utility>
42
43using namespace XrdClHttp;
44
45std::chrono::steady_clock::duration CurlOperation::m_stall_interval{CurlOperation::m_default_stall_interval};
47
48namespace {
49
50// For connection callbacks, we don't want to require a real DNS lookup; instead, we
51// will generate a fake address in the 169.254.x.y range and use that for the connection.
52// This will be fed to libcurl via the CURLOPT_RESOLVE option, which will bypass DNS lookups.
53
54// A randomized counter for generating fake addresses in the 169.254.x.y range
55thread_local int64_t fake_dns_counter = -1;
56
57// Map from hostname:port to fake address (e.g., 169.254.x.y:port) and
58// std::string pointer for the fake address.
59// We must track the std::string pointer so we can pass it to libcurl
60// as the CURLOPT_CLOSESOCKETDATA, which is passed to the close socket callback; the
61// lifetime of the pointer must be at least as long as the lifetime of the socket
62// (we maintain a reference count manually below).
63thread_local std::unordered_map<std::string, std::pair<std::string, std::string*>> fake_dns_map;
64
65// Reverse map from fake address (e.g., 169.254.x.y:port) to hostname:port and reference pointer
66thread_local std::unordered_map<std::string, std::pair<std::string, std::string*>> reverse_fake_dns_map;
67
68// References to fake addresses in use. The value is a reference count of sockets using
69// this address; when the count goes to zero, we can remove the entry from the above maps.
70// The second member is a unique_ptr to the std::string for the fake address, which will be
71// cleaned up when the refcount goes to zero.
72struct refcount_entry {
73 int count;
74 std::unique_ptr<std::string> addr;
75 std::chrono::steady_clock::time_point last_used;
76
77 bool IsExpired(std::chrono::steady_clock::time_point now) const {
78 return (now - last_used) > std::chrono::minutes(1);
79 }
80};
81
82thread_local std::unordered_map<std::string *, std::unique_ptr<refcount_entry>> fake_dns_refcount;
83
84std::string GenerateFakeEndpoint() {
85 if (fake_dns_counter == -1) {
86#ifdef __APPLE__
87 fake_dns_counter = arc4random();
88#else
89 errno = 0;
90 while (fake_dns_counter < 0 || errno == EINTR) {
91 if (getrandom((void*)&fake_dns_counter, sizeof(fake_dns_counter), 0) == sizeof(fake_dns_counter)) {
92 break;
93 }
94 }
95#endif
96 }
97 uint64_t addr = static_cast<uint64_t>(fake_dns_counter);
98 uint32_t class_d = addr & 0xff;
99 uint32_t class_c = (addr >> 8) & 0xff;
100 uint32_t port = 1024 + ((addr >> 16) % (65535 - 1024));
101 fake_dns_counter++;
102
103 return std::string("169.254.") + std::to_string(class_c) + "." + std::to_string(class_d) + ":" + std::to_string(port);
104}
105
106std::string *GetFakeEndpointForHost(const std::string &host, int port) {
107 std::string key = host + ":" + std::to_string(port);
108 auto it = fake_dns_map.find(key);
109 if (it != fake_dns_map.end()) {
110 return it->second.second;
111 }
112 auto addr = GenerateFakeEndpoint();
113 if (reverse_fake_dns_map.find(addr) != reverse_fake_dns_map.end()) {
114 return nullptr; // Collision, out of addresses.
115 }
116 auto addr_ptr_raw = new std::string(addr);
117 std::unique_ptr<std::string> addr_ptr(addr_ptr_raw);
118 fake_dns_map[key] = {addr, addr_ptr.get()};
119 reverse_fake_dns_map[addr] = {key, addr_ptr.get()};
120 std::unique_ptr<refcount_entry> new_entry(new refcount_entry{0, std::move(addr_ptr), std::chrono::steady_clock::now()});
121 fake_dns_refcount[addr_ptr_raw] = std::move(new_entry);
122 return addr_ptr_raw;
123}
124
125std::pair<std::string, int> ParseHostPort(const std::string &location) {
126 auto pos = location.find("://");
127 std::string authority = (pos == std::string::npos) ? location : location.substr(pos + 3);
128 std::string schema = (pos == std::string::npos) ? "" : location.substr(0, pos);
129 int std_port = (schema == "https" || schema == "davs") ? 443 : 80;
130 auto at_pos = authority.find('@');
131 std::string hostport = (at_pos == std::string::npos) ? authority : authority.substr(at_pos + 1);
132 pos = hostport.find('/');
133 if (pos != std::string::npos) {
134 hostport = hostport.substr(0, pos);
135 }
136 pos = hostport.find(':');
137 if (pos == std::string::npos) {
138 return {hostport, std_port};
139 }
140 int port = std_port;
141 try {
142 port = std::stoi(hostport.substr(pos + 1));
143 } catch (...) {
144 port = std_port;
145 }
146 return {hostport.substr(0, pos), port};
147}
148
149std::string DavToHttp(const std::string &url) {
150 if (url.compare(0, 6, "dav://") == 0) {
151 return "http://" + url.substr(6);
152 }
153 if (url.compare(0, 7, "davs://") == 0) {
154 return "https://" + url.substr(7);
155 }
156 return url;
157}
158
159} // namespace
160
161std::chrono::steady_clock::time_point CalculateExpiry(struct timespec timeout) {
162 if (timeout.tv_sec == 0 && timeout.tv_nsec == 0) {
163 return std::chrono::steady_clock::now() + std::chrono::seconds(30);
164 }
165 return std::chrono::steady_clock::now() + std::chrono::seconds(timeout.tv_sec) + std::chrono::nanoseconds(timeout.tv_nsec);
166}
167
169 struct timespec timeout, XrdCl::Log *logger, CreateConnCalloutType callout,
170 HeaderCallout *header_callout) :
171 CurlOperation::CurlOperation(handler, url, CalculateExpiry(timeout), logger, callout, header_callout)
172 {}
173
175 std::chrono::steady_clock::time_point expiry, XrdCl::Log *logger,
176 CreateConnCalloutType callout, HeaderCallout *header_callout) :
177 m_header_expiry(expiry),
178 m_header_callout(header_callout),
179 m_last_reset(std::chrono::steady_clock::now()),
180 m_last_header_reset(m_last_reset),
181 m_start_op(m_last_reset),
182 m_header_start(m_last_reset),
183 m_conn_callout(callout),
184 m_url(DavToHttp(url)),
186 m_handler(handler),
187 m_curl(nullptr, &curl_easy_cleanup),
188 m_logger(logger)
189 {}
190
192
193void
194CurlOperation::ExtendDeadline(struct timespec timeout)
195{
196 auto expiry = CalculateExpiry(timeout);
197 auto current = m_header_expiry.load(std::memory_order_relaxed);
198 while (expiry > current &&
199 !m_header_expiry.compare_exchange_weak(current, expiry,
200 std::memory_order_relaxed))
201 ;
202}
203
204void
205CurlOperation::Fail(uint16_t errCode, uint32_t errNum, const std::string &msg)
206{
207 SetDone(true);
208 if (m_handler == nullptr) {return;}
209 if (!msg.empty()) {
210 m_logger->Debug(kLogXrdClHttp, "curl operation failed with message: %s", msg.c_str());
211 } else {
212 m_logger->Debug(kLogXrdClHttp, "curl operation failed with status code %d", errNum);
213 }
214 auto status = new XrdCl::XRootDStatus(XrdCl::stError, errCode, errNum, msg);
215 auto handle = m_handler;
216 m_handler = nullptr;
217 handle->HandleResponse(status, nullptr);
218}
219
220int
221CurlOperation::FailCallback(XErrorCode ecode, const std::string &emsg) {
222 m_callback_error_code = ecode;
223 m_callback_error_str = emsg;
224 m_error = OpError::ErrCallback;
225 m_logger->Debug(kLogXrdClHttp, "%s", emsg.c_str());
226 return 0;
227}
228
229bool
231{
232 if(m_parsed_url)
233 {
235 }
236
237 // The client asks every request to carry the headers given by the HttpHeaders
238 // environment setting. HeaderBuilder interprets the contents.
239 std::string spec;
240 if (auto env = XrdCl::DefaultEnv::GetEnv()) {
241 env->GetString("HttpHeaders", spec);
242 }
243
244 if (!spec.empty()) {
245 // Refuse to send anything when the requested headers could not be understood;
246 // quietly omitting them changes what the request asks for.
247 //
248 // Fail here rather than leaving it to the caller, whose failure message for a
249 // setup error would suggest the wrong cause entirely.
250 HeaderList requested;
251 if (!HeaderBuilder::Build(spec, requested)) {
252 m_logger->Error(kLogXrdClHttp, "Not sending request to %s: the requested"
253 " headers could not be used", m_url.c_str());
254 Fail(XrdCl::errInvalidArgs, EINVAL, "Invalid header requested");
255 return false;
256 }
258 }
259
260 if (!m_header_callout) {
261 m_header_slist.reset();
262 for (const auto &header : m_headers_list) {
263 m_header_slist.reset(curl_slist_append(m_header_slist.release(),
264 (header.first + ": " + header.second).c_str()));
265 }
266 return curl_easy_setopt(curl, CURLOPT_HTTPHEADER, m_header_slist.get()) == CURLE_OK;
267 }
268 const auto &verb = GetVerbString(GetVerb());
269
270 auto extra_headers = m_header_callout->GetHeaders(
272 if (!extra_headers) {
274 "Failed to get headers from header callout for %s",
275 m_request_url.c_str());
276 return false;
277 }
278 m_header_slist.reset();
279 for (const auto &header : *extra_headers) {
280 if (HeaderBuilder::CompareIgnoreCase(header.first, "Content-Length")) {
281 auto upload_size = std::stoull(header.second);
282 curl_easy_setopt(curl, CURLOPT_INFILESIZE_LARGE, upload_size);
283 continue;
284 }
285 m_header_slist.reset(curl_slist_append(m_header_slist.release(),
286 (header.first + ": " + header.second).c_str()));
287 }
288 return curl_easy_setopt(curl, CURLOPT_HTTPHEADER, m_header_slist.get()) == CURLE_OK;
289}
290
291const std::string
293{
294 switch (verb) {
295 case HttpVerb::COPY:
296 return "COPY";
297 case HttpVerb::DELETE:
298 return "DELETE";
299 case HttpVerb::GET:
300 return "GET";
301 case HttpVerb::POST:
302 return "POST";
303 case HttpVerb::HEAD:
304 return "HEAD";
305 case HttpVerb::MKCOL:
306 return "MKCOL";
308 return "OPTIONS";
310 return "PROPFIND";
311 case HttpVerb::PUT:
312 return "PUT";
313 case HttpVerb::Count:
314 return "UNKNOWN";
315 }
316 return "UNKNOWN";
317}
318
319size_t
320CurlOperation::HeaderCallback(char *buffer, size_t size, size_t nitems, void *this_ptr)
321{
322 std::string header(buffer, size * nitems);
323 auto me = static_cast<CurlOperation*>(this_ptr);
324 auto now = std::chrono::steady_clock::now();
325 if (!me->m_received_header) {
326 me->m_received_header = true;
327 me->m_header_start = now;
328 }
329 me->m_header_lastop = now;
330 auto rv = me->Header(header);
331 return rv ? (size * nitems) : 0;
332}
333
334bool
335CurlOperation::Header(const std::string &header)
336{
337 auto result = m_headers.Parse(header);
338 // m_logger->Debug(kLogXrdClHttp, "Got header: %s", header.c_str());
339 if (!result) {
340 m_logger->Debug(kLogXrdClHttp, "Failed to parse response header: %s", header.c_str());
341 }
342 if (m_headers.HeadersDone()) {
343 if (!m_response_info) {
344 m_response_info.reset(new ResponseInfo());
345 }
346 m_response_info->AddResponse(m_headers.MoveHeaders());
347 }
348 return result;
349}
350
352CurlOperation::Redirect(std::string &target)
353{
354 m_callout.reset();
355 m_conn_callout_result = -1;
356 m_conn_callout_listener = -1;
357 m_tried_broker = false;
358
359 auto location = m_headers.GetLocation();
360 if (location.empty()) {
361 m_logger->Warning(kLogXrdClHttp, "After request to %s, server returned a redirect with no new location", m_request_url.c_str());
362 Fail(XrdCl::errErrorResponse, kXR_ServerError, "Server returned redirect without updated location");
364 }
365 if (location.size() && location[0] == '/') { // hostname not included in the location - redirect to self.
366 std::string_view orig_url(m_request_url);
367 auto scheme_loc = orig_url.find("://");
368 if (scheme_loc == std::string_view::npos) {
369 Fail(XrdCl::errErrorResponse, kXR_ServerError, "Server returned a location with unknown hostname");
371 }
372 auto path_loc = orig_url.find('/', scheme_loc + 3);
373 if (path_loc == std::string_view::npos) {
374 location = m_request_url + location;
375 } else {
376 location = std::string(orig_url.substr(0, path_loc)) + location;
377 }
378 }
379 m_logger->Debug(kLogXrdClHttp, "Request for %s redirected to %s", m_request_url.c_str(), location.c_str());
380 m_request_url = DavToHttp(location);
381 target = m_request_url;
382 curl_easy_setopt(m_curl.get(), CURLOPT_URL, m_request_url.c_str());
383 int disable_x509;
384 auto env = XrdCl::DefaultEnv::GetEnv();
385 if (env->GetInt("HttpDisableX509", disable_x509) && !disable_x509) {
386 std::string cert, key;
387 env->GetString("HttpClientCertFile", cert);
388 env->GetString("HttpClientKeyFile", key);
389 if (!cert.empty())
390 curl_easy_setopt(m_curl.get(), CURLOPT_SSLCERT, cert.c_str());
391 if (!key.empty())
392 curl_easy_setopt(m_curl.get(), CURLOPT_SSLKEY, key.c_str());
393 }
395
396 if (m_conn_callout) {
397 auto conn_callout = m_conn_callout(m_request_url, *m_response_info);
398 if (conn_callout != nullptr) {
399
400 auto [host, port] = ParseHostPort(m_request_url);
401 if (host.empty() || port == -1) {
402 Fail(XrdCl::errInternal, 0, "Failed to parse host and port from URL " + m_request_url);
404 }
405 auto fake_addr = GetFakeEndpointForHost(host, port);
406 if (!fake_addr || fake_addr->empty()) {
407 Fail(XrdCl::errInternal, 0, "Failed to generate a fake address for host " + host);
409 }
410 m_resolve_slist.reset(curl_slist_append(m_resolve_slist.release(),
411 (host + ":" + std::to_string(port) + ":" + *fake_addr).c_str()));
412 m_logger->Debug(kLogXrdClHttp, "For connection callout in redirect, mapping %s:%d -> %s", host.c_str(), port, fake_addr->c_str());
413
414 m_callout.reset(conn_callout);
415 std::string err;
417 // BeginCallout takes the expiration by non-const reference; hand it
418 // a copy rather than the atomic's storage.
419 auto expiry = GetHeaderExpiry();
420 if ((m_conn_callout_listener = m_callout->BeginCallout(err, expiry)) == -1) {
421 auto errMsg = "Failed to start a connection callout request: " + err;
422 Fail(XrdCl::errInternal, 0, errMsg.c_str());
424 }
425 curl_easy_setopt(m_curl.get(), CURLOPT_OPENSOCKETFUNCTION, CurlOperation::OpenSocketCallback);
426 curl_easy_setopt(m_curl.get(), CURLOPT_CLOSESOCKETFUNCTION, CurlOperation::CloseSocketCallback);
427 curl_easy_setopt(m_curl.get(), CURLOPT_OPENSOCKETDATA, this);
428 curl_easy_setopt(m_curl.get(), CURLOPT_CLOSESOCKETDATA, fake_addr);
429 curl_easy_setopt(m_curl.get(), CURLOPT_SOCKOPTFUNCTION, CurlOperation::SockOptCallback);
430 curl_easy_setopt(m_curl.get(), CURLOPT_SOCKOPTDATA, this);
431 curl_easy_setopt(m_curl.get(), CURLOPT_CONNECT_TO, m_resolve_slist.get());
432 }
433 }
434 m_received_header = false;
435
436 m_last_header_reset = m_last_reset = m_header_start = m_start_op = m_header_lastop = std::chrono::steady_clock::now();
438}
439
440namespace {
441
442size_t
443NullCallback(char * /*buffer*/, size_t size, size_t nitems, void * /*this_ptr*/)
444{
445 return size * nitems;
446}
447
448}
449
450void
452 m_is_paused = paused;
453 if (m_is_paused) {
454 m_pause_start = std::chrono::steady_clock::now();
455 } else if (m_pause_start != std::chrono::steady_clock::time_point{}) {
456 m_pause_duration += std::chrono::steady_clock::now() - m_pause_start;
457 m_pause_start = std::chrono::steady_clock::time_point{};
458 }
459}
460
461bool
463{
464 auto expiry = GetHeaderExpiry();
465 if ((m_conn_callout_listener = m_callout->BeginCallout(err, expiry)) == -1) {
466 err = "Failed to start a callout for a socket connection: " + err;
467 Fail(XrdCl::errInternal, 1, err.c_str());
468 return false;
469 }
470 return true;
471}
472
473std::tuple<uint64_t, std::chrono::steady_clock::duration, std::chrono::steady_clock::duration, std::chrono::steady_clock::duration>
475 auto now = std::chrono::steady_clock::now();
476 std::chrono::steady_clock::duration pre_header{}, post_header{}, pause_duration{};
477 if (m_received_header) {
478 if (m_last_header_reset < m_header_start) {
479 pre_header = m_header_start - m_last_header_reset;
480 m_last_header_reset = m_header_start;
481 }
482 post_header = now - ((m_last_reset < m_header_start) ? m_header_start : m_last_reset);
483 m_last_reset = now;
484 } else {
485 pre_header = now - m_last_header_reset;
486 m_last_header_reset = now;
487 }
488 if (IsPaused()) {
489 m_pause_duration += now - m_pause_start;
490 m_pause_start = now;
491 }
492 if (m_pause_duration != std::chrono::steady_clock::duration::zero()) {
493 pause_duration = m_pause_duration;
494 m_pause_duration = std::chrono::steady_clock::duration::zero();
495 }
496 auto bytes = m_bytes;
497 m_bytes = 0;
498 return {bytes, pre_header, post_header, pause_duration};
499}
500
501bool
502CurlOperation::HeaderTimeoutExpired(const std::chrono::steady_clock::time_point &now) {
503 if (m_received_header) return false;
504
505 if (now > GetHeaderExpiry()) {
506 if (m_error == OpError::ErrNone) m_error = OpError::ErrHeaderTimeout;
507 return true;
508 }
509 return false;
510}
511
512bool
513CurlOperation::OperationTimeoutExpired(const std::chrono::steady_clock::time_point &now) {
514 if (m_operation_expiry == std::chrono::steady_clock::time_point{} ||
515 !m_received_header) {
516 return false;
517 }
518
519 if (now > m_operation_expiry) {
520 if (m_error == OpError::ErrNone) m_error = OpError::ErrOperationTimeout;
521 return true;
522 }
523 return false;
524}
525
526bool
527CurlOperation::TransferStalled(uint64_t xfer, const std::chrono::steady_clock::time_point &now)
528{
529 // First, check to see how long it's been since any data was sent.
530 if (m_last_xfer == std::chrono::steady_clock::time_point()) {
531 m_last_xfer = m_header_lastop;
532 }
533 auto elapsed = now - m_last_xfer;
534 uint64_t xfer_diff = 0;
535 if (xfer > m_last_xfer_count) {
536 xfer_diff = xfer - m_last_xfer_count;
537 m_last_xfer_count = xfer;
538 m_last_xfer = now;
539 }
540
541 // If progress is made in this callback do not classify as stalled
542 if (elapsed > m_stall_interval && xfer_diff == 0) {
544 return true;
545 }
546
547 // Curl updated us with new timing but the byte count hasn't changed; no need to update the EMA.
548 if (xfer_diff == 0) {
549 return false;
550 }
551
552 // If the transfer is not stalled, then we check to see if the exponentially-weighted
553 // moving average of the transfer rate is below the minimum.
554
555 // If the stall interval since the last header hasn't passed, then we don't check for slow transfers.
556 auto elapsed_since_last_headerop = now - m_header_lastop;
557 if (elapsed_since_last_headerop < m_stall_interval) {
558 return false;
559 } else if (m_ema_rate < 0) {
560 m_ema_rate = xfer / std::chrono::duration<double>(elapsed_since_last_headerop).count();
561 }
562 // Calculate the exponential moving average of the transfer rate.
563 double elapsed_seconds = std::chrono::duration<double>(elapsed).count();
564 auto recent_rate = static_cast<double>(xfer_diff) / elapsed_seconds;
565 auto alpha = 1.0 - exp(-elapsed_seconds / std::chrono::duration<double>(m_stall_interval).count());
566 m_ema_rate = (1.0 - alpha) * m_ema_rate + alpha * recent_rate;
567 if (m_ema_rate < static_cast<double>(m_minimum_rate)) {
568 if (m_error == OpError::ErrNone) m_error = OpError::ErrTransferSlow;
569 return true;
570 }
571 return false;
572}
573
574bool
576{
577 if (curl == nullptr) {
578 throw std::runtime_error("Unable to setup curl operation with no handle");
579 }
580 struct timespec now;
581 if (clock_gettime(CLOCK_MONOTONIC, &now) == -1) {
582 throw std::runtime_error("Unable to get current time");
583 }
584
585 m_pause_start = {};
586 m_last_header_reset = m_last_reset = m_start_op = m_header_start = m_header_lastop = std::chrono::steady_clock::now();
587
588 m_curl.reset(curl);
589 m_curl_error_buffer[0] = '\0';
590 curl_easy_setopt(m_curl.get(), CURLOPT_URL, m_request_url.c_str());
591 curl_easy_setopt(m_curl.get(), CURLOPT_ERRORBUFFER, m_curl_error_buffer);
592 curl_easy_setopt(m_curl.get(), CURLOPT_HEADERFUNCTION, CurlStatOp::HeaderCallback);
593 curl_easy_setopt(m_curl.get(), CURLOPT_HEADERDATA, this);
594 curl_easy_setopt(m_curl.get(), CURLOPT_WRITEFUNCTION, NullCallback);
595 curl_easy_setopt(m_curl.get(), CURLOPT_WRITEDATA, nullptr);
596 curl_easy_setopt(m_curl.get(), CURLOPT_XFERINFOFUNCTION, CurlOperation::XferInfoCallback);
597 curl_easy_setopt(m_curl.get(), CURLOPT_XFERINFODATA, this);
598 curl_easy_setopt(m_curl.get(), CURLOPT_NOPROGRESS, 0L);
599 // Note: libcurl is not threadsafe unless this option is set.
600 // Before we set it, we saw deadlocks (and partial deadlocks) in practice.
601 curl_easy_setopt(m_curl.get(), CURLOPT_NOSIGNAL, 1L);
602
603 m_parsed_url = std::make_unique<XrdCl::URL>(m_request_url);
604 auto env = XrdCl::DefaultEnv::GetEnv();
605 int disable_x509;
606 if (env->GetInt("HttpDisableX509", disable_x509) && !disable_x509) {
607 auto [cert, key] = worker.ClientX509CertKeyFile();
608 if (!cert.empty()) {
610 "Using client X.509 credential found at %s", cert.c_str());
611 curl_easy_setopt(m_curl.get(), CURLOPT_SSLCERT, cert.c_str());
612 if (key.empty()) {
614 "X.509 client credential specified but not the client key");
615 } else {
616 curl_easy_setopt(m_curl.get(), CURLOPT_SSLKEY, key.c_str());
617 }
618 }
619 }
620
621 if (m_conn_callout) {
622 ResponseInfo info;
623 auto callout = m_conn_callout(m_request_url, info);
624 if (callout) {
625 m_callout.reset(callout);
626 m_conn_callout_listener = -1;
627 m_conn_callout_result = -1;
628 m_tried_broker = false;
629
630 auto [host, port] = ParseHostPort(m_request_url);
631 if (host.empty() || port == -1) {
632 throw std::runtime_error(
633 "Failed to parse host and port from URL " + m_request_url);
634 }
635 auto fake_addr = GetFakeEndpointForHost(host, port);
636 if (!fake_addr || fake_addr->empty()) {
637 throw std::runtime_error("Failed to generate a fake address for host " + host);
638 }
639 m_resolve_slist.reset(curl_slist_append(m_resolve_slist.release(),
640 (host + ":" + std::to_string(port) + ":" + *fake_addr).c_str()));
641 m_logger->Debug(kLogXrdClHttp, "For connection callout in operation setup, mapping %s:%d -> %s", host.c_str(), port, fake_addr->c_str());
642
643 curl_easy_setopt(m_curl.get(), CURLOPT_CONNECT_TO, m_resolve_slist.get());
644
645 curl_easy_setopt(m_curl.get(), CURLOPT_OPENSOCKETFUNCTION, CurlOperation::OpenSocketCallback);
646 curl_easy_setopt(m_curl.get(), CURLOPT_CLOSESOCKETFUNCTION, CurlOperation::CloseSocketCallback);
647 curl_easy_setopt(m_curl.get(), CURLOPT_OPENSOCKETDATA, this);
648 curl_easy_setopt(m_curl.get(), CURLOPT_CLOSESOCKETDATA, fake_addr);
649 curl_easy_setopt(m_curl.get(), CURLOPT_SOCKOPTFUNCTION, CurlOperation::SockOptCallback);
650 curl_easy_setopt(m_curl.get(), CURLOPT_SOCKOPTDATA, this);
651 }
652 }
653
654 return true;
655}
656
657bool
658CurlOperation::SetupNextRequest(const std::string &url, CurlWorker &worker)
659{
660 if (!m_curl) return false;
661
662 curl_easy_reset(m_curl.get());
663 ConfigureHandle(m_curl.get(), false);
664
665 m_request_url = DavToHttp(url);
667 m_headers_list.clear();
668 m_header_slist.reset();
669 m_response_info.reset();
670 m_resolve_slist.reset();
671 m_callout.reset();
672 m_conn_callout_listener = -1;
673 m_conn_callout_result = -1;
674 m_tried_broker = false;
675 m_received_header = false;
676 m_error = OpError::ErrNone;
677 m_callback_error_code = kXR_noErrorYet;
678 m_callback_error_str.clear();
679 m_last_xfer = {};
680 m_last_xfer_count = 0;
681 m_ema_rate = -1.0;
682
683 CURL *curl = m_curl.release();
684 return CurlOperation::Setup(curl, worker);
685}
686
687void
689{
690 m_conn_callout_listener = -1;
691 m_conn_callout_result = -1;
692 m_tried_broker = false;
693 m_callout.reset();
694
695 if (m_curl == nullptr) return;
696 curl_easy_setopt(m_curl.get(), CURLOPT_OPENSOCKETFUNCTION, nullptr);
697 curl_easy_setopt(m_curl.get(), CURLOPT_CLOSESOCKETFUNCTION, nullptr);
698 curl_easy_setopt(m_curl.get(), CURLOPT_OPENSOCKETDATA, nullptr);
699 curl_easy_setopt(m_curl.get(), CURLOPT_CLOSESOCKETDATA, nullptr);
700 curl_easy_setopt(m_curl.get(), CURLOPT_SOCKOPTFUNCTION, nullptr);
701 curl_easy_setopt(m_curl.get(), CURLOPT_SOCKOPTDATA, nullptr);
702 curl_easy_setopt(m_curl.get(), CURLOPT_SSLCERT, nullptr);
703 curl_easy_setopt(m_curl.get(), CURLOPT_SSLKEY, nullptr);
704 curl_easy_setopt(m_curl.get(), CURLOPT_HTTPHEADER, nullptr);
705 curl_easy_setopt(m_curl.get(), CURLOPT_CONNECT_TO, nullptr);
706 m_header_slist.reset();
707 m_curl.release();
708}
709
710curl_socket_t
711CurlOperation::OpenSocketCallback(void *clientp, curlsocktype purpose, struct curl_sockaddr *address)
712{
713 auto me = reinterpret_cast<CurlOperation*>(clientp);
714 auto fd = me->m_conn_callout_result;
715 me->m_conn_callout_result = -1;
716 if (fd == -1) {
717 std::string err;
718 auto expiry = me->GetHeaderExpiry();
719 if ((me->m_conn_callout_listener = me->m_callout->BeginCallout(err, expiry)) == -1) {
720 me->m_logger->Debug(kLogXrdClHttp, "Failed to start a connection callout request: %s", err.c_str());
721 }
722 return CURL_SOCKET_BAD;
723 } else {
724 sockaddr_in *inaddr = reinterpret_cast<sockaddr_in*>(&address->addr);
725 char ip_str[INET_ADDRSTRLEN];
726 char full_address_str[INET_ADDRSTRLEN + 6];
727 inet_ntop(AF_INET, &(inaddr->sin_addr), ip_str, INET_ADDRSTRLEN);
728 int port = ntohs(inaddr->sin_port);
729 snprintf(full_address_str, sizeof(full_address_str), "%s:%d", ip_str, port);
730 me->m_logger->Debug(kLogXrdClHttp, "Recording socket %d for %s", fd, full_address_str);
731 auto reverse_iter = reverse_fake_dns_map.find(full_address_str);
732 if (reverse_iter == reverse_fake_dns_map.end()) {
733 me->m_logger->Error(kLogXrdClHttp, "Failed to find fake DNS reverse entry for %s", full_address_str);
734 close(fd);
735 return CURL_SOCKET_BAD;
736 } else {
737 auto iter = fake_dns_refcount.find(reverse_iter->second.second);
738 if (iter == fake_dns_refcount.end()) {
739 me->m_logger->Error(kLogXrdClHttp, "Failed to find fake DNS refcount entry for %s", full_address_str);
740 close(fd);
741 return CURL_SOCKET_BAD;
742 }
743 iter->second->count++;
744 iter->second->last_used = std::chrono::steady_clock::now();
745 }
746
747 return fd;
748 }
749}
750
751int
752CurlOperation::SockOptCallback(void *clientp, curl_socket_t curlfd, curlsocktype purpose)
753{
754 return CURL_SOCKOPT_ALREADY_CONNECTED;
755}
756
757curl_socket_t
758CurlOperation::CloseSocketCallback(void *clientp, curl_socket_t fd)
759{
760 close(fd);
761 auto me = reinterpret_cast<std::string*>(clientp);
762 if (me == nullptr) {return 0;}
763 auto iter = fake_dns_refcount.find(me);
764 if (iter != fake_dns_refcount.end()) {
765 iter->second->count--;
766 if (iter->second->count <= 0 && iter->second->IsExpired(std::chrono::steady_clock::now())) {
767 auto rev_iter = reverse_fake_dns_map.find(*me);
768 if (rev_iter != reverse_fake_dns_map.end()) {
769 fake_dns_map.erase(rev_iter->second.first);
770 reverse_fake_dns_map.erase(rev_iter);
771 }
772 fake_dns_refcount.erase(iter);
773 }
774 }
775
776 return 0;
777}
778
779void
781{
782 auto now = std::chrono::steady_clock::now();
783 for (auto it = fake_dns_refcount.begin(); it != fake_dns_refcount.end(); ) {
784 if (it->second->count <= 0 && it->second->IsExpired(now)) {
785 auto rev_iter = reverse_fake_dns_map.find(*it->first);
786 if (rev_iter != reverse_fake_dns_map.end()) {
787 fake_dns_map.erase(rev_iter->second.first);
788 reverse_fake_dns_map.erase(rev_iter);
789 }
790 it = fake_dns_refcount.erase(it);
791 } else {
792 ++it;
793 }
794 }
795}
796
797int
798CurlOperation::XferInfoCallback(void *clientp, curl_off_t /*dltotal*/, curl_off_t dlnow, curl_off_t /*ultotal*/, curl_off_t ulnow)
799{
800 auto me = reinterpret_cast<CurlOperation*>(clientp);
801 auto now = std::chrono::steady_clock::now();
802 if (me->HeaderTimeoutExpired(now) || me->OperationTimeoutExpired(now)) {
803 return 1; // return value triggers CURLE_ABORTED_BY_CALLBACK
804 }
805 uint64_t xfer_bytes = dlnow > ulnow ? dlnow : ulnow;
806 if (me->TransferStalled(xfer_bytes, now)) {
807 return 1;
808 }
809 return 0;
810}
811
812int
814{
815 m_conn_callout_result = m_callout ? m_callout->FinishCallout(err) : -1;
816 if (m_callout && m_conn_callout_result == -1) {
817 m_logger->Error(kLogXrdClHttp, "Error when getting socket callout: %s", err.c_str());
818 } else if (m_callout) {
819 m_logger->Debug(kLogXrdClHttp, "Got callback socket %d", m_conn_callout_result);
820 }
821 return m_conn_callout_result;
822}
XErrorCode
@ kXR_noErrorYet
@ kXR_ServerError
std::chrono::steady_clock::time_point CalculateExpiry(struct timespec timeout)
void CURL
#define close(a)
Definition XrdPosix.hh:48
int emsg(int rc, char *msg)
void SetDone(bool has_failed)
int FailCallback(XErrorCode ecode, const std::string &emsg)
std::chrono::steady_clock::time_point GetHeaderExpiry() const
bool FinishSetup(CURL *curl)
std::atomic< std::chrono::steady_clock::time_point > m_header_expiry
const std::string m_url
std::unique_ptr< CURL, void(*)(CURL *)> m_curl
bool TransferStalled(uint64_t xfer_bytes, const std::chrono::steady_clock::time_point &now)
static const std::string GetVerbString(HttpVerb)
virtual HttpVerb GetVerb() const =0
virtual void ReleaseHandle()
static void CleanupDnsCache()
std::tuple< uint64_t, std::chrono::steady_clock::duration, std::chrono::steady_clock::duration, std::chrono::steady_clock::duration > StatisticsReset()
bool SetupNextRequest(const std::string &url, CurlWorker &worker)
static constexpr int m_default_minimum_rate
std::vector< std::pair< std::string, std::string > > HeaderList
std::vector< std::pair< std::string, std::string > > m_headers_list
HeaderCallout * m_header_callout
bool HeaderTimeoutExpired(const std::chrono::steady_clock::time_point &now)
virtual int WaitSocketCallback(std::string &err)
void ExtendDeadline(struct timespec timeout)
XrdClHttp::HttpVerb HttpVerb
std::chrono::steady_clock::time_point m_operation_expiry
virtual void Fail(uint16_t errCode, uint32_t errNum, const std::string &)
virtual RedirectAction Redirect(std::string &target)
XrdCl::ResponseHandler * m_handler
CurlOperation(XrdCl::ResponseHandler *handler, const std::string &url, struct timespec timeout, XrdCl::Log *log, CreateConnCalloutType, HeaderCallout *header_callout)
void SetPaused(bool paused)
bool StartConnectionCallout(std::string &err)
bool OperationTimeoutExpired(const std::chrono::steady_clock::time_point &now)
virtual bool Setup(CURL *curl, CurlWorker &)
std::tuple< std::string, std::string > ClientX509CertKeyFile() const
virtual std::shared_ptr< HeaderList > GetHeaders(const std::string &verb, const std::string &url, const HeaderList &headers)=0
ResponseInfo::HeaderMap && MoveHeaders()
bool Parse(const std::string &headers)
const std::string & GetLocation() const
static Env * GetEnv()
Get default client environment.
Handle diagnostics.
Definition XrdClLog.hh:101
void Error(uint64_t topic, const char *format,...)
Report an error.
Definition XrdClLog.cc:231
void Warning(uint64_t topic, const char *format,...)
Report a warning.
Definition XrdClLog.cc:248
void Debug(uint64_t topic, const char *format,...)
Print a debug message.
Definition XrdClLog.cc:282
Handle an async response.
virtual void HandleResponse(XRootDStatus *status, AnyObject *response)
bool CompareIgnoreCase(const std::string_view lhs, const std::string_view rhs)
void AppendMissing(const HeaderList &extra, HeaderList &headers)
bool Build(const std::string_view spec, HeaderList &headers)
ConnectionCallout *(*)(const std::string &, const ResponseInfo &) CreateConnCalloutType
void InjectBearerToken(const XrdCl::URL &url, std::vector< std::pair< std::string, std::string > > &headers, XrdCl::Log *logger=nullptr)
const uint64_t kLogXrdClHttp
void ConfigureHandle(CURL *curl, bool verbose)
const uint16_t errErrorResponse
const uint16_t stError
An error occurred that could potentially be retried.
const uint16_t errInternal
Internal error.
const uint16_t errInvalidArgs