39#include <sys/random.h>
45std::chrono::steady_clock::duration CurlOperation::m_stall_interval{CurlOperation::m_default_stall_interval};
55thread_local int64_t fake_dns_counter = -1;
63thread_local std::unordered_map<std::string, std::pair<std::string, std::string*>> fake_dns_map;
66thread_local std::unordered_map<std::string, std::pair<std::string, std::string*>> reverse_fake_dns_map;
72struct refcount_entry {
74 std::unique_ptr<std::string> addr;
75 std::chrono::steady_clock::time_point last_used;
77 bool IsExpired(std::chrono::steady_clock::time_point now)
const {
78 return (now - last_used) > std::chrono::minutes(1);
82thread_local std::unordered_map<std::string *, std::unique_ptr<refcount_entry>> fake_dns_refcount;
84std::string GenerateFakeEndpoint() {
85 if (fake_dns_counter == -1) {
87 fake_dns_counter = arc4random();
90 while (fake_dns_counter < 0 || errno == EINTR) {
91 if (getrandom((
void*)&fake_dns_counter,
sizeof(fake_dns_counter), 0) ==
sizeof(fake_dns_counter)) {
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));
103 return std::string(
"169.254.") + std::to_string(class_c) +
"." + std::to_string(class_d) +
":" + std::to_string(port);
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;
112 auto addr = GenerateFakeEndpoint();
113 if (reverse_fake_dns_map.find(addr) != reverse_fake_dns_map.end()) {
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);
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);
136 pos = hostport.find(
':');
137 if (pos == std::string::npos) {
138 return {hostport, std_port};
142 port = std::stoi(hostport.substr(pos + 1));
146 return {hostport.substr(0, pos), port};
149std::string DavToHttp(
const std::string &url) {
150 if (url.compare(0, 6,
"dav://") == 0) {
151 return "http://" + url.substr(6);
153 if (url.compare(0, 7,
"davs://") == 0) {
154 return "https://" + url.substr(7);
162 if (timeout.tv_sec == 0 && timeout.tv_nsec == 0) {
163 return std::chrono::steady_clock::now() + std::chrono::seconds(30);
165 return std::chrono::steady_clock::now() + std::chrono::seconds(timeout.tv_sec) + std::chrono::nanoseconds(timeout.tv_nsec);
175 std::chrono::steady_clock::time_point expiry,
XrdCl::Log *logger,
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)),
187 m_curl(nullptr, &curl_easy_cleanup),
198 while (expiry > current &&
200 std::memory_order_relaxed))
222 m_callback_error_code = ecode;
223 m_callback_error_str =
emsg;
241 env->GetString(
"HttpHeaders", spec);
253 " headers could not be used",
m_url.c_str());
261 m_header_slist.reset();
263 m_header_slist.reset(curl_slist_append(m_header_slist.release(),
264 (header.first +
": " + header.second).c_str()));
266 return curl_easy_setopt(curl, CURLOPT_HTTPHEADER, m_header_slist.get()) == CURLE_OK;
272 if (!extra_headers) {
274 "Failed to get headers from header callout for %s",
278 m_header_slist.reset();
279 for (
const auto &header : *extra_headers) {
281 auto upload_size = std::stoull(header.second);
282 curl_easy_setopt(curl, CURLOPT_INFILESIZE_LARGE, upload_size);
285 m_header_slist.reset(curl_slist_append(m_header_slist.release(),
286 (header.first +
": " + header.second).c_str()));
288 return curl_easy_setopt(curl, CURLOPT_HTTPHEADER, m_header_slist.get()) == CURLE_OK;
320CurlOperation::HeaderCallback(
char *buffer,
size_t size,
size_t nitems,
void *this_ptr)
322 std::string header(buffer, size * nitems);
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;
329 me->m_header_lastop = now;
330 auto rv = me->Header(header);
331 return rv ? (size * nitems) : 0;
335CurlOperation::Header(
const std::string &header)
343 if (!m_response_info) {
344 m_response_info.reset(
new ResponseInfo());
355 m_conn_callout_result = -1;
356 m_conn_callout_listener = -1;
357 m_tried_broker =
false;
360 if (location.empty()) {
365 if (location.size() && location[0] ==
'/') {
367 auto scheme_loc = orig_url.find(
"://");
368 if (scheme_loc == std::string_view::npos) {
372 auto path_loc = orig_url.find(
'/', scheme_loc + 3);
373 if (path_loc == std::string_view::npos) {
376 location = std::string(orig_url.substr(0, path_loc)) + location;
385 if (env->GetInt(
"HttpDisableX509", disable_x509) && !disable_x509) {
386 std::string cert, key;
387 env->GetString(
"HttpClientCertFile", cert);
388 env->GetString(
"HttpClientKeyFile", key);
390 curl_easy_setopt(
m_curl.get(), CURLOPT_SSLCERT, cert.c_str());
392 curl_easy_setopt(
m_curl.get(), CURLOPT_SSLKEY, key.c_str());
396 if (m_conn_callout) {
397 auto conn_callout = m_conn_callout(
m_request_url, *m_response_info);
398 if (conn_callout !=
nullptr) {
401 if (host.empty() || port == -1) {
405 auto fake_addr = GetFakeEndpointForHost(host, port);
406 if (!fake_addr || fake_addr->empty()) {
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());
414 m_callout.reset(conn_callout);
420 if ((m_conn_callout_listener = m_callout->BeginCallout(err, expiry)) == -1) {
421 auto errMsg =
"Failed to start a connection callout request: " + err;
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());
434 m_received_header =
false;
436 m_last_header_reset = m_last_reset = m_header_start = m_start_op = m_header_lastop = std::chrono::steady_clock::now();
443NullCallback(
char * ,
size_t size,
size_t nitems,
void * )
445 return size * nitems;
452 m_is_paused = 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{};
465 if ((m_conn_callout_listener = m_callout->BeginCallout(err, expiry)) == -1) {
466 err =
"Failed to start a callout for a socket connection: " + err;
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;
482 post_header = now - ((m_last_reset < m_header_start) ? m_header_start : m_last_reset);
485 pre_header = now - m_last_header_reset;
486 m_last_header_reset = now;
489 m_pause_duration += now - m_pause_start;
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();
496 auto bytes = m_bytes;
498 return {bytes, pre_header, post_header, pause_duration};
503 if (m_received_header)
return false;
515 !m_received_header) {
530 if (m_last_xfer == std::chrono::steady_clock::time_point()) {
531 m_last_xfer = m_header_lastop;
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;
542 if (elapsed > m_stall_interval && xfer_diff == 0) {
548 if (xfer_diff == 0) {
556 auto elapsed_since_last_headerop = now - m_header_lastop;
557 if (elapsed_since_last_headerop < m_stall_interval) {
559 }
else if (m_ema_rate < 0) {
560 m_ema_rate = xfer / std::chrono::duration<double>(elapsed_since_last_headerop).count();
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;
577 if (curl ==
nullptr) {
578 throw std::runtime_error(
"Unable to setup curl operation with no handle");
581 if (clock_gettime(CLOCK_MONOTONIC, &now) == -1) {
582 throw std::runtime_error(
"Unable to get current time");
586 m_last_header_reset = m_last_reset = m_start_op = m_header_start = m_header_lastop = std::chrono::steady_clock::now();
589 m_curl_error_buffer[0] =
'\0';
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);
601 curl_easy_setopt(
m_curl.get(), CURLOPT_NOSIGNAL, 1L);
606 if (env->GetInt(
"HttpDisableX509", disable_x509) && !disable_x509) {
610 "Using client X.509 credential found at %s", cert.c_str());
611 curl_easy_setopt(
m_curl.get(), CURLOPT_SSLCERT, cert.c_str());
614 "X.509 client credential specified but not the client key");
616 curl_easy_setopt(
m_curl.get(), CURLOPT_SSLKEY, key.c_str());
621 if (m_conn_callout) {
625 m_callout.reset(callout);
626 m_conn_callout_listener = -1;
627 m_conn_callout_result = -1;
628 m_tried_broker =
false;
631 if (host.empty() || port == -1) {
632 throw std::runtime_error(
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);
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());
643 curl_easy_setopt(
m_curl.get(), CURLOPT_CONNECT_TO, m_resolve_slist.get());
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);
660 if (!
m_curl)
return false;
662 curl_easy_reset(
m_curl.get());
668 m_header_slist.reset();
669 m_response_info.reset();
670 m_resolve_slist.reset();
672 m_conn_callout_listener = -1;
673 m_conn_callout_result = -1;
674 m_tried_broker =
false;
675 m_received_header =
false;
678 m_callback_error_str.clear();
680 m_last_xfer_count = 0;
690 m_conn_callout_listener = -1;
691 m_conn_callout_result = -1;
692 m_tried_broker =
false;
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();
711CurlOperation::OpenSocketCallback(
void *clientp, curlsocktype purpose,
struct curl_sockaddr *address)
714 auto fd = me->m_conn_callout_result;
715 me->m_conn_callout_result = -1;
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());
722 return CURL_SOCKET_BAD;
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);
735 return CURL_SOCKET_BAD;
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);
741 return CURL_SOCKET_BAD;
743 iter->second->count++;
744 iter->second->last_used = std::chrono::steady_clock::now();
752CurlOperation::SockOptCallback(
void *clientp, curl_socket_t curlfd, curlsocktype purpose)
754 return CURL_SOCKOPT_ALREADY_CONNECTED;
758CurlOperation::CloseSocketCallback(
void *clientp, curl_socket_t 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);
772 fake_dns_refcount.erase(iter);
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);
790 it = fake_dns_refcount.erase(it);
798CurlOperation::XferInfoCallback(
void *clientp, curl_off_t , curl_off_t dlnow, curl_off_t , curl_off_t ulnow)
801 auto now = std::chrono::steady_clock::now();
802 if (me->HeaderTimeoutExpired(now) || me->OperationTimeoutExpired(now)) {
805 uint64_t xfer_bytes = dlnow > ulnow ? dlnow : ulnow;
806 if (me->TransferStalled(xfer_bytes, now)) {
815 m_conn_callout_result = m_callout ? m_callout->FinishCallout(err) : -1;
816 if (m_callout && m_conn_callout_result == -1) {
818 }
else if (m_callout) {
821 return m_conn_callout_result;
std::chrono::steady_clock::time_point CalculateExpiry(struct timespec timeout)
int emsg(int rc, char *msg)
void SetDone(bool has_failed)
int FailCallback(XErrorCode ecode, const std::string &emsg)
static int m_minimum_transfer_rate
std::chrono::steady_clock::time_point GetHeaderExpiry() const
bool FinishSetup(CURL *curl)
std::atomic< std::chrono::steady_clock::time_point > m_header_expiry
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::string m_request_url
std::tuple< std::string, std::string > ClientX509CertKeyFile() const
static Env * GetEnv()
Get default client environment.
void Error(uint64_t topic, const char *format,...)
Report an error.
void Warning(uint64_t topic, const char *format,...)
Report a warning.
void Debug(uint64_t topic, const char *format,...)
Print a debug message.
Handle an async response.
virtual void HandleResponse(XRootDStatus *status, AnyObject *response)
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