XRootD
Loading...
Searching...
No Matches
XrdClHttpUtil.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 "XrdClHttpFile.hh"
22#include "XrdClHttpOps.hh"
24#include "XrdClHttpUtil.hh"
25#include "XrdClHttpWorker.hh"
26
29#include <XrdCl/XrdClLog.hh>
30#include <XrdCl/XrdClURL.hh>
31#include <XrdCl/XrdClUtils.hh>
33#include <XrdCks/XrdCksData.hh>
34#include <XrdOuc/XrdOucCRC.hh>
37#include <XrdVersion.hh>
38
39#include <curl/curl.h>
40#include <openssl/bio.h>
41#include <openssl/evp.h>
42
43#include <fcntl.h>
44#include <fstream>
45#ifdef __APPLE__
46#include <pthread.h>
47#else
48#include <sys/syscall.h>
49#include <sys/types.h>
50#endif
51#include <unistd.h>
52
53#include <algorithm>
54#include <cerrno>
55#include <cctype>
56#include <cstdlib>
57#include <sstream>
58#include <stdexcept>
59#include <system_error>
60#include <utility>
61
62using namespace XrdClHttp;
63
64thread_local std::vector<CURL*> HandlerQueue::m_handles;
65std::atomic<unsigned> CurlWorker::m_maintenance_period = 5;
66std::vector<std::unique_ptr<XrdClHttp::CurlWorker>> CurlWorker::m_workers;
67std::mutex CurlWorker::m_workers_mutex;
68
69// Performance statistics for the worker
70std::atomic<uint64_t> CurlWorker::m_conncall_errors = 0;
71std::atomic<uint64_t> CurlWorker::m_conncall_req = 0;
72std::atomic<uint64_t> CurlWorker::m_conncall_success = 0;
73std::atomic<uint64_t> CurlWorker::m_conncall_timeout = 0;
74decltype(CurlWorker::m_ops) CurlWorker::m_ops = {};
75std::vector<std::atomic<std::chrono::system_clock::rep>*> CurlWorker::m_workers_last_completed_cycle;
76std::vector<std::atomic<std::chrono::system_clock::rep>*> CurlWorker::m_workers_oldest_op;
77std::mutex CurlWorker::m_worker_stats_mutex;
78
79// Performance statistics for the queue
80std::atomic<uint64_t> HandlerQueue::m_ops_consumed = 0; // Count of operations consumed from the queue.
81std::atomic<uint64_t> HandlerQueue::m_ops_produced = 0; // Count of operations added to the queue.
82std::atomic<uint64_t> HandlerQueue::m_ops_rejected = 0; // Count of operations rejected by the queue.
83
84// shutdown + init trigger, must be last of the static members
85CurlWorker::initcontrol CurlWorker::m_initcontrol;
86
88 CURL *curl{nullptr};
89 time_t expiry{0};
90};
91
92namespace {
93
94pid_t getthreadid() {
95#if defined(__APPLE__)
96 uint64_t pth_threadid;
97 pthread_threadid_np(pthread_self(), &pth_threadid);
98 return pth_threadid;
99#elif defined(__linux__)
100 // NOTE: glibc 2.30 finally provides a gettid() wrapper; however,
101 // we currently support RHEL 8, which is based on glibc 2.28. Until
102 // we drop that platform, it's easier to do the syscall directly on Linux
103 // instead of additional ifdef calls.
104 return syscall(SYS_gettid);
105#else
106 return getpid();
107#endif
108}
109
110}
111
112bool XrdClHttp::HTTPStatusIsError(unsigned status) {
113 return (status < 100) || (status >= 400);
114}
115
117{
118 auto env = XrdCl::DefaultEnv::GetEnv();
119 if (!env) return {};
120
121 std::string token;
122 if (!env->GetString("BearerToken", token) || token.empty()) {
123 env->ImportString("BearerToken", "BEARER_TOKEN");
124 env->GetString("BearerToken", token);
125 }
126 if (!token.empty()) {
127 XrdCl::Utils::Trim(token);
128 return token;
129 }
130
131 std::string token_file;
132 if (!env->GetString("BearerTokenFile", token_file) || token_file.empty()) {
133 env->ImportString("BearerTokenFile", "BEARER_TOKEN_FILE");
134 env->GetString("BearerTokenFile", token_file);
135 }
136 if (token_file.empty()) return {};
137
138 std::ifstream input(token_file);
139 if (!input) return {};
140 std::getline(input, token);
141 XrdCl::Utils::Trim(token);
142 return token;
143}
144
145bool XrdClHttp::ShouldUseBearerToken(const std::string &protocols,
146 bool hasX509Credential,
147 bool hasBearerToken)
148{
149 if(protocols.empty()) return hasBearerToken;
150
151 std::vector<std::string> requested;
152 XrdCl::Utils::splitString(requested, protocols, ",");
153 for(auto protocol : requested)
154 {
155 XrdCl::Utils::Trim(protocol);
156 std::transform(protocol.begin(), protocol.end(), protocol.begin(),
157 [](unsigned char c) { return std::tolower(c); });
158 if(protocol == "gsi" && hasX509Credential) return false;
159 if(protocol == "ztn" && hasBearerToken) return true;
160 }
161 return false;
162}
163
165 const XrdCl::URL &url,
166 std::vector<std::pair<std::string, std::string>> &headers,
167 XrdCl::Log *logger)
168{
169 if(url.GetParams().count("authz")) return;
170 for(const auto &header : headers)
171 {
172 if(header.first == "Authorization") return;
173 }
174
175 const std::string token = GetBearerToken(logger);
176 if(token.empty()) return;
177
178 bool hasX509Credential = false;
179 auto env = XrdCl::DefaultEnv::GetEnv();
180 if(env)
181 {
182 int disableX509 = 0;
183 std::string certificate;
184 hasX509Credential =
185 env->GetInt("HttpDisableX509", disableX509)
186 && !disableX509
187 && env->GetString("HttpClientCertFile", certificate)
188 && !certificate.empty();
189 }
190
191 const char *configuredProtocols = std::getenv("XrdSecPROTOCOL");
192 if(!ShouldUseBearerToken(configuredProtocols ? configuredProtocols : "",
193 hasX509Credential, true))
194 {
195 return;
196 }
197
198 if(logger)
199 {
200 logger->Debug(kLogXrdClHttp,
201 "Injecting bearer token from environment for %s",
202 url.GetURL().c_str());
203 }
204 headers.emplace_back("Authorization", "Bearer " + token);
205}
206
207std::pair<uint16_t, uint32_t> XrdClHttp::HTTPStatusConvert(unsigned status) {
208 switch (status) {
209 case 400: // Bad Request
210 return std::make_pair(XrdCl::errErrorResponse, kXR_InvalidRequest);
211 case 401: // Unauthorized (needs authentication)
212 return std::make_pair(XrdCl::errErrorResponse, kXR_NotAuthorized);
213 case 402: // Payment Required
214 case 403: // Forbidden (failed authorization)
215 return std::make_pair(XrdCl::errErrorResponse, kXR_NotAuthorized);
216 case 404:
217 return std::make_pair(XrdCl::errErrorResponse, kXR_NotFound);
218 case 405: // Method not allowed
219 return std::make_pair(XrdCl::errErrorResponse, kXR_Unsupported);
220 case 406: // Not acceptable
221 return std::make_pair(XrdCl::errErrorResponse, kXR_InvalidRequest);
222 case 407: // Proxy Authentication Required
223 return std::make_pair(XrdCl::errErrorResponse, kXR_NotAuthorized);
224 case 408: // Request timeout
225 return std::make_pair(XrdCl::errErrorResponse, kXR_ReqTimedOut);
226 case 409: // Conflict
227 return std::make_pair(XrdCl::errErrorResponse, kXR_Conflict);
228 case 410: // Gone
229 return std::make_pair(XrdCl::errErrorResponse, kXR_NotFound);
230 case 411: // Length required
231 case 412: // Precondition failed
232 case 413: // Payload too large
233 case 414: // URI too long
234 case 415: // Unsupported Media Type
235 case 416: // Range Not Satisfiable
236 case 417: // Expectation Failed
237 case 418: // I'm a teapot
238 return std::make_pair(XrdCl::errErrorResponse, kXR_InvalidRequest);
239 case 421: // Misdirected Request
240 case 422: // Unprocessable Content
241 return std::make_pair(XrdCl::errErrorResponse, kXR_InvalidRequest);
242 case 423: // Locked
243 return std::make_pair(XrdCl::errErrorResponse, kXR_FileLocked);
244 case 424: // Failed Dependency
245 case 425: // Too Early
246 case 426: // Upgrade Required
247 case 428: // Precondition Required
248 return std::make_pair(XrdCl::errErrorResponse, kXR_InvalidRequest);
249 case 429: // Too Many Requests
250 return std::make_pair(XrdCl::errErrorResponse, kXR_Overloaded);
251 case 431: // Request Header Fields Too Large
252 return std::make_pair(XrdCl::errErrorResponse, kXR_InvalidRequest);
253 case 451: // Unavailable For Legal Reasons
254 return std::make_pair(XrdCl::errErrorResponse, kXR_Impossible);
255 case 500: // Internal Server Error
256 return std::make_pair(XrdCl::errErrorResponse, kXR_ServerError);
257 case 501: // Not Implemented
258 return std::make_pair(XrdCl::errErrorResponse, kXR_Unsupported);
259 case 502: // Bad Gateway
260 return std::make_pair(XrdCl::errErrorResponse, kXR_ServerError);
261 case 503: // Service Unavailable
262 return std::make_pair(XrdCl::errErrorResponse, kXR_Overloaded);
263 case 504: // Gateway Timeout
264 return std::make_pair(XrdCl::errErrorResponse, kXR_ReqTimedOut);
265 case 507: // Insufficient Storage
266 return std::make_pair(XrdCl::errErrorResponse, kXR_overQuota);
267 case 508: // Loop Detected
268 case 510: // Not Extended
269 case 511: // Network Authentication Required
270 return std::make_pair(XrdCl::errErrorResponse, kXR_ServerError);
271 }
272 return std::make_pair(XrdCl::errUnknown, status);
273}
274
275std::pair<uint16_t, uint32_t> CurlCodeConvert(CURLcode res) {
276 switch (res) {
277 case CURLE_OK:
278 return std::make_pair(XrdCl::errNone, 0);
279 case CURLE_COULDNT_RESOLVE_PROXY:
280 case CURLE_COULDNT_RESOLVE_HOST:
281 return std::make_pair(XrdCl::errInvalidAddr, 0);
282 case CURLE_LOGIN_DENIED:
283 // Commented-out cases are for platforms (RHEL7) where the error
284 // codes are undefined.
285 //case CURLE_AUTH_ERROR:
286 //case CURLE_SSL_CLIENTCERT:
287 case CURLE_REMOTE_ACCESS_DENIED:
288 return std::make_pair(XrdCl::errLoginFailed, EACCES);
289 case CURLE_SSL_CONNECT_ERROR:
290 case CURLE_SSL_ENGINE_NOTFOUND:
291 case CURLE_SSL_ENGINE_SETFAILED:
292 case CURLE_SSL_CERTPROBLEM:
293 case CURLE_SSL_CIPHER:
294 case 51: // In old curl versions, this is CURLE_PEER_FAILED_VERIFICATION; that constant was changed to be 60 / CURLE_SSL_CACERT
295 case CURLE_SSL_SHUTDOWN_FAILED:
296 case CURLE_SSL_CRL_BADFILE:
297 case CURLE_SSL_ISSUER_ERROR:
298 case CURLE_SSL_CACERT: // value is 60; merged with CURLE_PEER_FAILED_VERIFICATION
299 //case CURLE_SSL_PINNEDPUBKEYNOTMATCH:
300 //case CURLE_SSL_INVALIDCERTSTATUS:
301 return std::make_pair(XrdCl::errTlsError, 0);
302 case CURLE_SEND_ERROR:
303 case CURLE_RECV_ERROR:
304 return std::make_pair(XrdCl::errSocketError, EIO);
305 case CURLE_COULDNT_CONNECT:
306 case CURLE_GOT_NOTHING:
307 return std::make_pair(XrdCl::errConnectionError, ECONNREFUSED);
308 case CURLE_OPERATION_TIMEDOUT:
309#ifdef HAVE_XPROTOCOL_TIMEREXPIRED
311#else
312 return std::make_pair(XrdCl::errOperationExpired, ESTALE);
313#endif
314 case CURLE_UNSUPPORTED_PROTOCOL:
315 case CURLE_NOT_BUILT_IN:
316 return std::make_pair(XrdCl::errNotSupported, ENOSYS);
317 case CURLE_FAILED_INIT:
318 return std::make_pair(XrdCl::errInternal, 0);
319 case CURLE_URL_MALFORMAT:
320 return std::make_pair(XrdCl::errInvalidArgs, res);
321 //case CURLE_WEIRD_SERVER_REPLY:
322 //case CURLE_HTTP2:
323 //case CURLE_HTTP2_STREAM:
324 return std::make_pair(XrdCl::errCorruptedHeader, res);
325 case CURLE_PARTIAL_FILE:
326 return std::make_pair(XrdCl::errDataError, res);
327 // These two errors indicate a failure in the callback. That
328 // should generate their own failures, meaning this should never
329 // get use.
330 case CURLE_READ_ERROR:
331 case CURLE_WRITE_ERROR:
332 return std::make_pair(XrdCl::errInternal, res);
333 case CURLE_RANGE_ERROR:
334 case CURLE_BAD_CONTENT_ENCODING:
335 return std::make_pair(XrdCl::errNotSupported, res);
336 case CURLE_TOO_MANY_REDIRECTS:
337 return std::make_pair(XrdCl::errRedirectLimit, res);
338 default:
339 return std::make_pair(XrdCl::errUnknown, res);
340 }
341}
342
344 std::string_view input,
345 std::array<unsigned char, g_max_checksum_length> &output) {
346 if (input.size() > 4 * ((output.size() + 2) / 3)
347 || input.size() % 4 != 0)
348 return false;
349 if (input.size() == 0) return true;
350
351 std::unique_ptr<BIO, decltype(&BIO_free_all)> b64(BIO_new(BIO_f_base64()), &BIO_free_all);
352 BIO_set_flags(b64.get(), BIO_FLAGS_BASE64_NO_NL);
353 std::unique_ptr<BIO, decltype(&BIO_free_all)> bmem(
354 BIO_new_mem_buf(const_cast<char *>(input.data()), input.size()), &BIO_free_all);
355 bmem.reset(BIO_push(b64.release(), bmem.release()));
356
357 // Compute expected length of output; used to verify BIO_read consumes all input
358 size_t expectedLen = static_cast<size_t>(input.size() * 0.75);
359 if (input[input.size() - 1] == '=') {
360 expectedLen -= 1;
361 if (input[input.size() - 2] == '=') {
362 expectedLen -= 1;
363 }
364 }
365
366 auto len = BIO_read(bmem.get(), &output[0], output.size());
367
368 if (len == -1 || static_cast<size_t>(len) != expectedLen) return false;
369
370 return true;
371}
372
373// Parse a single header line.
374//
375// Curl promises for its callbacks "The header callback is
376// called once for each header and only complete header lines
377// are passed on to the callback".
378bool HeaderParser::Parse(const std::string &header_line)
379{
380 if (m_recv_all_headers) {
381 m_recv_all_headers = false;
382 m_recv_status_line = false;
383 }
384
385 if (!m_recv_status_line) {
386 m_recv_status_line = true;
387
388 std::stringstream ss(header_line);
389 std::string item;
390 if (!std::getline(ss, item, ' ')) return false;
391 m_resp_protocol = item;
392 if (!std::getline(ss, item, ' ')) return false;
393 try {
394 m_status_code = std::stol(item);
395 } catch (...) {
396 return false;
397 }
398 if (m_status_code < 100 || m_status_code >= 600) {
399 return false;
400 }
401 if (!std::getline(ss, item, '\n')) return false;
402 auto cr_loc = item.find('\r');
403 if (cr_loc != std::string::npos) {
404 m_resp_message = item.substr(0, cr_loc);
405 } else {
406 m_resp_message = item;
407 }
408 return true;
409 }
410
411 if (header_line.empty() || header_line == "\n" || header_line == "\r\n") {
412 m_recv_all_headers = true;
413 return true;
414 }
415
416 auto found = header_line.find(":");
417 if (found == std::string::npos) {
418 return false;
419 }
420
421 std::string header_name = header_line.substr(0, found);
422 if (!Canonicalize(header_name)) {
423 return false;
424 }
425
426 found += 1;
427 while (found < header_line.size()) {
428 if (header_line[found] != ' ') {break;}
429 found += 1;
430 }
431 std::string header_value = header_line.substr(found);
432 // Note: ignoring the fact headers are only supposed to contain ASCII.
433 // We should trim out UTF-8.
434 header_value.erase(header_value.find_last_not_of(" \r\n\t") + 1);
435
436 // Record the line in our header structure. Will be returned as part
437 // of the response info object.
438 auto iter = m_headers.find(header_name);
439 if (iter == m_headers.end()) {
440 m_headers.insert(iter, {header_name, {header_value}});
441 } else {
442 iter->second.push_back(header_value);
443 }
444
445 if (header_name == "Allow") {
446 std::string_view val(header_value);
447 while (!val.empty()) {
448 auto found = val.find(',');
449 auto method = val.substr(0, found);
450 if (method == "PROPFIND") {
451 auto new_verbs = static_cast<unsigned>(m_allow_verbs) | static_cast<unsigned>(VerbsCache::HttpVerb::kPROPFIND);
452 m_allow_verbs = static_cast<VerbsCache::HttpVerb>(new_verbs);
453 }
454 if (found == std::string_view::npos) break;
455 val = val.substr(found + 1);
456 }
457 if (static_cast<unsigned>(m_allow_verbs) & ~static_cast<unsigned>(VerbsCache::HttpVerb::kUnknown)) {
458 m_allow_verbs = static_cast<VerbsCache::HttpVerb>(static_cast<unsigned>(m_allow_verbs) & ~static_cast<unsigned>(VerbsCache::HttpVerb::kUnknown));
459 }
460 } else if (header_name == "Content-Length") {
461 try {
462 m_content_length = std::stoll(header_value);
463 } catch (...) {
464 return false;
465 }
466 }
467 else if (header_name == "Content-Type") {
468 std::string_view val(header_value);
469 auto found = val.find(";");
470 auto first_type = val.substr(0, found);
471 m_multipart_byteranges = first_type == "multipart/byteranges";
472 if (m_multipart_byteranges) {
473 auto remainder = val.substr(found + 1);
474 found = remainder.find("boundary=");
475 if (found != std::string_view::npos) {
476 SetMultipartSeparator(remainder.substr(found + 9));
477 }
478 }
479 }
480 else if (header_name == "Content-Range") {
481 auto found = header_value.find(" ");
482 if (found == std::string::npos) {
483 return false;
484 }
485 std::string range_unit = header_value.substr(0, found);
486 if (range_unit != "bytes") {
487 return false;
488 }
489 auto range_resp = header_value.substr(found + 1);
490 found = range_resp.find("/");
491 if (found == std::string::npos) {
492 return false;
493 }
494 auto incl_range = range_resp.substr(0, found);
495 found = incl_range.find("-");
496 if (found == std::string::npos) {
497 return false;
498 }
499 auto first_pos = incl_range.substr(0, found);
500 try {
501 m_response_offset = std::stoll(first_pos);
502 } catch (...) {
503 return false;
504 }
505 auto last_pos = incl_range.substr(found + 1);
506 size_t last_byte;
507 try {
508 last_byte = std::stoll(last_pos);
509 } catch (...) {
510 return false;
511 }
512 m_content_length = last_byte - m_response_offset + 1;
513 }
514 else if (header_name == "Location") {
515 m_location = header_value;
516 } else if (header_name == "Digest") {
517 ParseDigest(header_value, m_checksums);
518 }
519 else if (header_name == "Etag")
520 {
521 // Note, the original hader name is ETag, renamed to Etag in parsing
522 // remove additional quotes
523 m_etag = header_value;
524 m_etag.erase(remove(m_etag.begin(), m_etag.end(), '\"'), m_etag.end());
525 }
526 else if (header_name == "Cache-Control")
527 {
528 m_cache_control = header_value;
529 }
530
531 return true;
532}
533
534// Parse a RFC 3230 header into the checksum info structure
535//
536// If the parsing fails, the second element of the tuple will be false.
537void HeaderParser::ParseDigest(const std::string &digest, XrdClHttp::ChecksumInfo &info) {
538 std::string_view view(digest);
539 std::array<unsigned char, g_max_checksum_length> checksum_value{};
540 std::string digest_lower;
541 while (!view.empty()) {
542 auto nextsep = view.find(',');
543 auto entry = view.substr(0, nextsep);
544 if (nextsep == std::string_view::npos) {
545 view = "";
546 } else {
547 view = view.substr(nextsep + 1);
548 }
549 nextsep = entry.find('=');
550 if (nextsep == std::string_view::npos) continue;
551 auto name = trim_view(entry.substr(0, nextsep));
552 auto value = trim_view(entry.substr(nextsep + 1));
553 digest_lower.clear();
554 digest_lower.resize(name.size());
555 std::transform(name.begin(), name.end(), digest_lower.begin(), [](unsigned char c) {
556 return std::tolower(c);
557 });
558 const auto setHex32 = [&](ChecksumType type) {
559 if (value.empty() || value.size() > 8) return;
560 std::string padded(8 - value.size(), '0');
561 padded.append(value.data(), value.size());
562 XrdCksData checksum;
563 if (!checksum.Set(padded.c_str(), padded.size())) return;
564 std::copy_n(
565 reinterpret_cast<const unsigned char *>(checksum.Value),
566 GetChecksumLength(type), checksum_value.begin());
567 info.Set(type, checksum_value);
568 };
569 const auto setBase64 = [&](ChecksumType type) {
570 const size_t decodedSize = GetChecksumLength(type);
571 const size_t encodedSize = 4 * ((decodedSize + 2) / 3);
572 if (value.size() == encodedSize
573 && Base64Decode(value, checksum_value))
574 {
575 info.Set(type, checksum_value);
576 }
577 };
578 if (digest_lower == "adler" || digest_lower == "adler32") {
579 setHex32(ChecksumType::kADLER32);
580 } else if (digest_lower == "crc32") {
581 setHex32(ChecksumType::kCRC32);
582 } else if (digest_lower == "md5") {
583 setBase64(ChecksumType::kMD5);
584 } else if (digest_lower == "sha" || digest_lower == "sha1") {
585 setBase64(ChecksumType::kSHA1);
586 } else if (digest_lower == "sha-256" || digest_lower == "sha256") {
587 setBase64(ChecksumType::kSHA256);
588 } else if (digest_lower == "crc32c") {
589 // XRootD currently incorrectly base64-encodes crc32c checksums; see
590 // https://github.com/xrootd/xrootd/issues/2456
591 // For backward comaptibility, if this looks like base64 encoded (8
592 // bytes long and last two bytes are padding), then we base64 decode.
593 if (value.size() == 8 && value[6] == '=' && value[7] == '=') {
594 setBase64(ChecksumType::kCRC32C);
595 continue;
596 }
597 setHex32(ChecksumType::kCRC32C);
598 }
599 }
600}
601
602// Convert the checksum type to a RFC 3230 digest name as recorded by IANA here:
603// https://www.iana.org/assignments/http-dig-alg/http-dig-alg.xhtml
605 switch (type) {
607 return "adler32";
609 return "crc32";
611 return "MD5";
613 return "CRC32c";
615 return "SHA";
617 return "SHA-256";
619 std::string result;
620 for (int value = 0;
621 value < static_cast<int>(XrdClHttp::ChecksumType::kAll);
622 ++value) {
623 if (!result.empty()) result += ',';
624 result += ChecksumTypeToDigestName(
625 static_cast<XrdClHttp::ChecksumType>(value));
626 }
627 return result;
628 }
629 default:
630 return "";
631 }
632}
633
634// This clever approach was inspired by golang's net/textproto
635bool HeaderParser::validHeaderByte(unsigned char c)
636{
637 const static uint64_t mask_lower = 0 |
638 uint64_t((1<<10)-1) << '0' |
639 uint64_t(1) << '!' |
640 uint64_t(1) << '#' |
641 uint64_t(1) << '$' |
642 uint64_t(1) << '%' |
643 uint64_t(1) << '&' |
644 uint64_t(1) << '\'' |
645 uint64_t(1) << '*' |
646 uint64_t(1) << '+' |
647 uint64_t(1) << '-' |
648 uint64_t(1) << '.';
649
650 const static uint64_t mask_upper = 0 |
651 uint64_t((1<<26)-1) << ('a'-64) |
652 uint64_t((1<<26)-1) << ('A'-64) |
653 uint64_t(1) << ('^'-64) |
654 uint64_t(1) << ('_'-64) |
655 uint64_t(1) << ('`'-64) |
656 uint64_t(1) << ('|'-64) |
657 uint64_t(1) << ('~'-64);
658
659 if (c >= 128) return false;
660 if (c >= 64) return (uint64_t(1)<<(c-64)) & mask_upper;
661 return (uint64_t(1) << c) & mask_lower;
662}
663
664bool HeaderParser::Canonicalize(std::string &headerName)
665{
666 auto upper = true;
667 const static int toLower = 'a' - 'A';
668 for (size_t idx=0; idx<headerName.size(); idx++) {
669 char c = headerName[idx];
670 if (!validHeaderByte(c)) {
671 return false;
672 }
673 if (upper && 'a' <= c && c <= 'z') {
674 c -= toLower;
675 } else if (!upper && 'A' <= c && c <= 'Z') {
676 c += toLower;
677 }
678 headerName[idx] = c;
679 upper = c == '-';
680 }
681 return true;
682}
683
684HandlerQueue::HandlerQueue(unsigned max_pending_ops) :
685 m_max_pending_ops(max_pending_ops)
686{
687 int filedes[2];
688 auto result = pipe(filedes);
689 if (result == -1) {
690 throw std::system_error(errno, std::generic_category(),
691 "failed to create HTTP worker pipe");
692 }
693 if (fcntl(filedes[0], F_SETFL, O_NONBLOCK | O_CLOEXEC) == -1 || fcntl(filedes[1], F_SETFL, O_NONBLOCK | O_CLOEXEC) == -1) {
694 const int error = errno;
695 close(filedes[0]);
696 close(filedes[1]);
697 throw std::system_error(error, std::generic_category(),
698 "failed to configure HTTP worker pipe");
699 }
700 m_read_fd = filedes[0];
701 m_write_fd = filedes[1];
702};
703
704namespace {
705
706bool EnableCurlHeaderDump() {
707 auto *log = XrdCl::DefaultEnv::GetLog();
708 if (log && log->GetLevel() >= XrdCl::Log::DumpMsg)
709 return true;
710
711 return false;
712}
713
714// Debug callback for libcurl headers; enabled with XRD_LOGLEVEL=Dump
715int DumpHeader(CURL *handle, curl_infotype type, char *data, size_t size, void *clientp) {
716 (void)handle;
717 auto *logger = static_cast<XrdCl::Log *>(clientp);
718 if (!logger || !data || size == 0) {
719 return 0;
720 }
721
722 const char *direction = nullptr;
723 switch (type) {
724 case CURLINFO_HEADER_OUT:
725 direction = ">";
726 break;
727 case CURLINFO_HEADER_IN:
728 direction = "<";
729 break;
730 default:
731 return 0;
732 }
733
734 const std::string redacted = obfuscateAuth(std::string(data, size));
735 logger->Debug(kLogXrdClHttp, "%s %s", direction, redacted.c_str());
736 return 0;
737}
738
739}
740
741// Trim left and right side of a string_view for space characters
742std::string_view XrdClHttp::trim_view(const std::string_view &input_view) {
743 auto view = XrdClHttp::ltrim_view(input_view);
744 for (size_t idx = 0; idx < input_view.size(); idx++) {
745 if (!isspace(view[view.size() - 1 - idx])) {
746 return view.substr(0, view.size() - idx);
747 }
748 }
749 return "";
750}
751
752// Trim the left side of a string_view for space
753std::string_view XrdClHttp::ltrim_view(const std::string_view &input_view) {
754 for (size_t idx = 0; idx < input_view.size(); idx++) {
755 if (!isspace(input_view[idx])) {
756 return input_view.substr(idx);
757 }
758 }
759 return "";
760}
761
762void
763XrdClHttp::ConfigureHandle(CURL *curl, bool verbose) {
764 curl_easy_setopt(curl, CURLOPT_USERAGENT, "xrdcl-http/" XrdVERSION);
765 curl_easy_setopt(curl, CURLOPT_DEBUGFUNCTION, DumpHeader);
766 curl_easy_setopt(curl, CURLOPT_DEBUGDATA, XrdCl::DefaultEnv::GetLog());
767 if (verbose)
768 curl_easy_setopt(curl, CURLOPT_VERBOSE, 1L);
769
770 auto env = XrdCl::DefaultEnv::GetEnv();
771 std::string ca_file;
772 if (!env->GetString("HttpCertFile", ca_file) || ca_file.empty()) {
773 char *x509_ca_file = getenv("X509_CERT_FILE");
774 if (x509_ca_file) {
775 ca_file = std::string(x509_ca_file);
776 }
777 }
778 if (!ca_file.empty()) {
779 curl_easy_setopt(curl, CURLOPT_CAINFO, ca_file.c_str());
780 }
781 std::string ca_dir;
782 if (!env->GetString("HttpCertDir", ca_dir) || ca_dir.empty()) {
783 char *x509_ca_dir = getenv("X509_CERT_DIR");
784 if (x509_ca_dir) {
785 ca_dir = std::string(x509_ca_dir);
786 }
787 }
788 if (!ca_dir.empty()) {
789 curl_easy_setopt(curl, CURLOPT_CAPATH, ca_dir.c_str());
790 }
791
792 curl_easy_setopt(curl, CURLOPT_BUFFERSIZE, 32*1024);
793}
794
795CURL *
796XrdClHttp::GetHandle(bool verbose) {
797 auto result = curl_easy_init();
798 if (result == nullptr) {
799 return result;
800 }
801
802 XrdClHttp::ConfigureHandle(result, verbose);
803
804 return result;
805}
806
807CURL *
809 if (m_handles.size()) {
810 auto result = m_handles.back();
811 m_handles.pop_back();
812 return result;
813 }
814
815 return ::GetHandle(EnableCurlHeaderDump());
816}
817
818void
820 m_handles.push_back(curl);
821}
822
823void
825{
826 std::unique_lock<std::mutex> lk(m_mutex);
827 auto now = std::chrono::steady_clock::now();
828
829 // Iterate through the paused transfers, checking if they are done.
830 for (auto &op : m_ops) {
831 if (!op->IsPaused()) continue;
832
833 if (op->TransferStalled(0, now)) {
834 op->ContinueHandle();
835 }
836 }
837
838 std::vector<decltype(m_ops)::value_type> expired_ops;
839 unsigned expired_count = 0;
840 auto it = std::remove_if(m_ops.begin(), m_ops.end(),
841 [&](const std::shared_ptr<CurlOperation> &handler) {
842 auto expired = handler->GetOperationExpiry() < now;
843 if (expired) {
844 expired_ops.push_back(handler);
845 expired_count++;
846 }
847 return expired;
848 });
849 m_ops.erase(it, m_ops.end());
850
851 // The contents of our pipe and the in-memory queue are now off by expired_count.
852 // Read exactly that many bytes from the pipe and throw them away.
853 char throwaway[64];
854 unsigned bytes_to_read = expired_count;
855 while (bytes_to_read > 0) {
856 size_t chunk = std::min<size_t>(sizeof(throwaway), bytes_to_read);
857 ssize_t n = read(m_read_fd, throwaway, chunk);
858 if (n > 0) {
859 bytes_to_read -= n;
860 } else if (n == -1) {
861 if (errno == EINTR) {
862 continue;
863 } else {
864 // EWOULDBLOCK is a possibility if there's a synchronization error;
865 // for now, just continue on as if we were successful in reading out
866 // the missing bytes
867 break;
868 }
869 } else {
870 break;
871 }
872 }
873
874 // Note: the failure handler may trigger new operations submitted to the queue
875 // (which requires the lock to be held) such as a prefetch operation that gets split
876 // into multiple sub-operations.
877 //
878 // Thus, we must unlock the mutex protecting the queue and avoid touching the shared state of
879 // m_ops.
880 lk.unlock();
881 for (auto &handler : expired_ops) {
882 if (handler) handler->Fail(XrdCl::errOperationExpired, 0, "Operation expired while in queue");
883 }
884}
885
886void
887HandlerQueue::Produce(std::shared_ptr<CurlOperation> handler)
888{
889 auto handler_expiry = handler->GetOperationExpiry();
890 std::unique_lock<std::mutex> lk{m_mutex};
891 m_producer_cv.wait_until(lk,
892 handler_expiry,
893 [&]{return m_ops.size() < m_max_pending_ops;}
894 );
895 if (std::chrono::steady_clock::now() > handler_expiry) {
896 lk.unlock();
897 handler->Fail(XrdCl::errOperationExpired, 0, "Operation expired while waiting for worker");
898 m_ops_rejected.fetch_add(1, std::memory_order_relaxed);
899 return;
900 }
901
902 m_ops.push_back(handler);
903 char ready[] = "1";
904 while (true) {
905 auto result = write(m_write_fd, ready, 1);
906 if (result == -1) {
907 if (errno == EINTR) {
908 continue;
909 } else if (errno == EAGAIN || errno == EWOULDBLOCK) {
910 // This should never happen, but if it does, just continue
911 // as if we successfully wrote the notification to the pipe.
912 break;
913 }
914 throw std::system_error(errno, std::generic_category(),
915 "failed to notify HTTP worker");
916 }
917 break;
918 }
919
920 lk.unlock();
921 m_consumer_cv.notify_one();
922 m_ops_produced.fetch_add(1, std::memory_order_relaxed);
923}
924
925std::shared_ptr<CurlOperation>
926HandlerQueue::Consume(std::chrono::steady_clock::duration dur)
927{
928 std::unique_lock<std::mutex> lk(m_mutex);
929 m_consumer_cv.wait_for(lk, dur, [&]{return m_ops.size() > 0 || m_shutdown;});
930 if (m_shutdown || m_ops.empty()) {
931 return {};
932 }
933
934 std::shared_ptr<CurlOperation> result = m_ops.front();
935 m_ops.pop_front();
936
937 char ready[1];
938 while (true) {
939 auto result = read(m_read_fd, ready, 1);
940 if (result == -1) {
941 if (errno == EINTR) {
942 continue;
943 } else if (errno == EAGAIN || errno == EWOULDBLOCK) {
944 // This should never happen, but if it does, just continue
945 // as if we successfully read the byte.
946 break;
947 }
948 throw std::system_error(errno, std::generic_category(),
949 "failed to consume HTTP worker notification");
950 }
951 break;
952 }
953
954 lk.unlock();
955 m_producer_cv.notify_one();
956 m_ops_consumed.fetch_add(1, std::memory_order_relaxed);
957
958 return result;
959}
960
961std::string
963{
964 auto consumed = m_ops_consumed.load(std::memory_order_relaxed);
965 auto produced = m_ops_produced.load(std::memory_order_relaxed);
966 return "{"
967 "\"produced\":" + std::to_string(produced) + ","
968 "\"consumed\":" + std::to_string(consumed) + ","
969 "\"pending\":" + std::to_string(produced - consumed) + ","
970 "\"rejected\":" + std::to_string(m_ops_rejected.load(std::memory_order_relaxed)) +
971 "}";
972}
973
974std::shared_ptr<CurlOperation>
976{
977 std::unique_lock<std::mutex> lk(m_mutex);
978 if (m_ops.size() == 0) {
979 std::shared_ptr<CurlOperation> result;
980 return result;
981 }
982
983 std::shared_ptr<CurlOperation> result = m_ops.front();
984 m_ops.pop_front();
985
986 char ready[1];
987 while (true) {
988 auto result = read(m_read_fd, ready, 1);
989 if (result == -1) {
990 if (errno == EINTR) {
991 continue;
992 } else if (errno == EAGAIN || errno == EWOULDBLOCK) {
993 // This should never happen, but if it does, just continue
994 // as if we successfully read the byte.
995 break;
996 }
997 throw std::system_error(errno, std::generic_category(),
998 "failed to consume HTTP worker notification");
999 }
1000 break;
1001 }
1002
1003 lk.unlock();
1004 m_producer_cv.notify_one();
1005 m_ops_consumed.fetch_add(1, std::memory_order_relaxed);
1006
1007 return result;
1008}
1009
1010void
1012{
1013 std::unique_lock lock(m_mutex);
1014 m_shutdown = true;
1015 m_consumer_cv.notify_all();
1016}
1017
1018void
1020{
1021 for (auto handle : m_handles) {
1022 curl_easy_cleanup(handle);
1023 }
1024 m_handles.clear();
1025}
1026
1027CurlWorker::CurlWorker(std::shared_ptr<HandlerQueue> queue, VerbsCache &cache, XrdCl::Log* logger) :
1028 m_cache(cache),
1029 m_queue(queue),
1030 m_logger(logger)
1031{
1032 {
1033 std::unique_lock lk(m_worker_stats_mutex);
1034 m_stats_offset = m_workers_last_completed_cycle.size();
1035 m_workers_last_completed_cycle.push_back(&m_last_completed_cycle);
1036 m_workers_oldest_op.push_back(&m_oldest_op);
1037 }
1038 int pipeInfo[2];
1039 if ((pipe(pipeInfo) == -1) || (fcntl(pipeInfo[0], F_SETFD, FD_CLOEXEC)) || (fcntl(pipeInfo[1], F_SETFD, FD_CLOEXEC))) {
1040 throw std::runtime_error("Failed to create shutdown monitoring pipe for curl worker");
1041 }
1042 m_shutdown_pipe_r = pipeInfo[0];
1043 m_shutdown_pipe_w = pipeInfo[1];
1044
1045 // Handle setup of the X509 authentication
1046 auto env = XrdCl::DefaultEnv::GetEnv();
1047 env->GetString("HttpClientCertFile", m_x509_client_cert_file);
1048 env->GetString("HttpClientKeyFile", m_x509_client_key_file);
1049}
1050
1051std::tuple<std::string, std::string> CurlWorker::ClientX509CertKeyFile() const
1052{
1053 return std::make_tuple(m_x509_client_cert_file, m_x509_client_key_file);
1054}
1055
1056std::string
1058{
1059 auto now = std::chrono::system_clock::now().time_since_epoch().count();
1060 auto oldest_op = now;
1061 auto oldest_cycle = now;
1062 {
1063 std::unique_lock lk(m_worker_stats_mutex);
1064 for (const auto &entry : m_workers_last_completed_cycle) {
1065 if (!entry) {continue;}
1066 auto cycle = entry->load(std::memory_order_relaxed);
1067 if (cycle < oldest_cycle) oldest_cycle = cycle;
1068 }
1069 for (const auto &entry : m_workers_oldest_op) {
1070 if (!entry) {continue;}
1071 auto op = entry->load(std::memory_order_relaxed);
1072 if (op < oldest_op) oldest_op = op;
1073 }
1074 }
1075 auto oldest_op_dbl = std::chrono::duration<double>(std::chrono::system_clock::time_point(std::chrono::system_clock::duration(oldest_op)).time_since_epoch()).count();
1076 auto oldest_cycle_dbl = std::chrono::duration<double>(std::chrono::system_clock::time_point(std::chrono::system_clock::duration(oldest_cycle)).time_since_epoch()).count();
1077 std::string retval = "{"
1078 "\"oldest_op\":" + std::to_string(oldest_op_dbl) + ","
1079 "\"oldest_cycle\":" + std::to_string(oldest_cycle_dbl) + ","
1080 ;
1081
1082 for (size_t verb_idx = 0; verb_idx < static_cast<int>(XrdClHttp::CurlOperation::HttpVerb::Count); verb_idx++) {
1083 const auto &verb_str = XrdClHttp::CurlOperation::GetVerbString(static_cast<XrdClHttp::CurlOperation::HttpVerb>(verb_idx));
1084 for (size_t op_idx = 0; op_idx < 402; op_idx++) {
1085 if (op_idx == 401) continue;
1086
1087 auto &op_stats = m_ops[verb_idx][op_idx];
1088 auto duration = op_stats.m_duration.load(std::memory_order_relaxed);
1089 if (duration == 0) continue;
1090
1091 std::string prefix = "http_" + verb_str + "_" + ((op_idx == 402) ? "invalid" : std::to_string(200 + op_idx)) + "_";
1092
1093 auto duration_dbl = std::chrono::duration<double>(std::chrono::steady_clock::duration(duration)).count();
1094 retval += "\"" + prefix + "duration\":" + std::to_string(duration_dbl) + ",";
1095
1096 duration = op_stats.m_pause_duration.load(std::memory_order_relaxed);
1097 if (duration > 0) {
1098 duration_dbl = std::chrono::duration<double>(std::chrono::steady_clock::duration(duration)).count();
1099 retval += "\"" + prefix + "pause_duration\":" + std::to_string(duration_dbl) + ",";
1100 }
1101
1102 auto count = op_stats.m_bytes.load(std::memory_order_relaxed);
1103 if (count) retval += "\"" + prefix + "bytes\":" + std::to_string(count) + ",";
1104 count = op_stats.m_error.load(std::memory_order_relaxed);
1105 if (count) retval += "\"" + prefix + "error\":" + std::to_string(count) + ",";
1106 count = op_stats.m_finished.load(std::memory_order_relaxed);
1107 if (count) retval += "\"" + prefix + "finished\":" + std::to_string(count) + ",";
1108 count = op_stats.m_client_timeout.load(std::memory_order_relaxed);
1109 if (count) retval += "\"" + prefix + "client_timeout\":" + std::to_string(count) + ",";
1110 count = op_stats.m_server_timeout.load(std::memory_order_relaxed);
1111 if (count) retval += "\"" + prefix + "server_timeout\":" + std::to_string(count) + ",";
1112 }
1113 {
1114 auto &op_stats = m_ops[verb_idx][401];
1115 auto duration = op_stats.m_duration.load(std::memory_order_relaxed);
1116 if (duration == 0) continue;
1117
1118 std::string prefix = "http_" + verb_str + "_";
1119
1120 auto duration_dbl = std::chrono::duration<double>(std::chrono::steady_clock::duration(duration)).count();
1121 retval += "\"" + prefix + "preheader_duration\":" + std::to_string(duration_dbl) + ",";
1122
1123 auto count = op_stats.m_started.load(std::memory_order_relaxed);
1124 if (count) retval += "\"" + prefix + "started\":" + std::to_string(count) + ",";
1125 count = op_stats.m_error.load(std::memory_order_relaxed);
1126 if (count) retval += "\"" + prefix + "preheader_error\":" + std::to_string(count) + ",";
1127 count = op_stats.m_finished.load(std::memory_order_relaxed);
1128 if (count) retval += "\"" + prefix + "preheader_finished\":" + std::to_string(count) + ",";
1129 count = op_stats.m_server_timeout.load(std::memory_order_relaxed);
1130 if (count) retval += "\"" + prefix + "preheader_timeout\":" + std::to_string(count) + ",";
1131 count = op_stats.m_conncall_timeout.load(std::memory_order_relaxed);
1132 if (count) retval += "\"" + prefix + "conncall_timeout\":" + std::to_string(count) + ",";
1133 }
1134 }
1135
1136 retval +=
1137 "\"conncall_error\":" + std::to_string(m_conncall_errors.load(std::memory_order_relaxed)) + ","
1138 "\"conncall_started\":" + std::to_string(m_conncall_req.load(std::memory_order_relaxed)) + ","
1139 "\"conncall_success\":" + std::to_string(m_conncall_success.load(std::memory_order_relaxed)) + ","
1140 "\"conncall_timeout\":" + std::to_string(m_conncall_timeout.load(std::memory_order_relaxed)) +
1141 "}";
1142
1143 return retval;
1144}
1145
1146void
1147CurlWorker::OpRecord(XrdClHttp::CurlOperation &op, OpKind kind)
1148{
1149 int sc = op.GetStatusCode();
1150 // - We encode everything pre-header as integer "401". We include a 100-continue request as "pre-header".
1151 // - Status codes out of the acceptable range are labeled "402"
1152 // - Otherwise, we store it in the array shifted by 200 (to avoid more sparsity)
1153 if (sc < 0 || kind == OpKind::Start || sc == 100) {
1154 sc = 401;
1155 } else if (sc < 200 || sc >= 600) {
1156 sc = 402;
1157 } else {
1158 sc -= 200;
1159 }
1160 auto [bytes, pre_headers, post_headers, pause_duration] = op.StatisticsReset();
1161 auto &op_stats = m_ops[static_cast<int>(op.GetVerb())][sc];
1162 op_stats.m_bytes.fetch_add(bytes, std::memory_order_relaxed);
1163 op_stats.m_duration.fetch_add((sc == 401) ? pre_headers.count() : post_headers.count(), std::memory_order_relaxed);
1164 op_stats.m_pause_duration.fetch_add(pause_duration.count(), std::memory_order_relaxed);
1165 if (pre_headers != std::chrono::steady_clock::duration::zero() && sc != 401) {
1166 auto &old_stats = m_ops[static_cast<int>(op.GetVerb())][401];
1167 old_stats.m_duration.fetch_add(pre_headers.count(), std::memory_order_relaxed);
1168 }
1169 switch (kind) {
1170 case OpKind::ConncallTimeout:
1171 op_stats.m_conncall_timeout.fetch_add(1, std::memory_order_relaxed);
1172 break;
1173 case OpKind::ClientTimeout:
1174 op_stats.m_client_timeout.fetch_add(1, std::memory_order_relaxed);
1175 break;
1176 case OpKind::Error:
1177 op_stats.m_error.fetch_add(1, std::memory_order_relaxed);
1178 break;
1179 case OpKind::Finish:
1180 op_stats.m_finished.fetch_add(1, std::memory_order_relaxed);
1181 break;
1182 case OpKind::Start:
1183 op_stats.m_started.fetch_add(1, std::memory_order_relaxed);
1184 break;
1185 case OpKind::ServerTimeout:
1186 op_stats.m_server_timeout.fetch_add(1, std::memory_order_relaxed);
1187 break;
1188 case OpKind::Update:
1189 break;
1190 }
1191}
1192
1193void
1194CurlWorker::Start(std::unique_ptr<XrdClHttp::CurlWorker> self, std::thread tid)
1195{
1196 {
1197 std::unique_lock lock(m_workers_mutex);
1198 m_workers.emplace_back(std::move(self));
1199 m_self_tid = std::move(tid);
1200 }
1201 std::unique_lock lock(m_start_lock);
1202 m_start_complete = true;
1203 m_start_complete_cv.notify_one();
1204}
1205
1206void
1208{
1209 {
1210 std::unique_lock lock(myself->m_start_lock);
1211 myself->m_start_complete_cv.wait(lock, [&]{return myself->m_start_complete;});
1212 }
1213 try {
1214 myself->Run();
1215 } catch (...) {
1216 myself->m_logger->Warning(kLogXrdClHttp, "Curl worker got an exception");
1217 {
1218 std::unique_lock lock(m_workers_mutex);
1219 auto iter = std::remove_if(m_workers.begin(), m_workers.end(), [&](std::unique_ptr<XrdClHttp::CurlWorker> &worker){return worker.get() == myself;});
1220 m_workers.erase(iter);
1221 }
1222 }
1223}
1224
1225void
1227 int max_pending = 50;
1228 XrdCl::DefaultEnv::GetEnv()->GetInt("HttpMaxPendingOps", max_pending);
1229 m_continue_queue.reset(new HandlerQueue(max_pending));
1230 auto &queue = *m_queue.get();
1231 m_logger->Debug(kLogXrdClHttp, "Started a curl worker");
1232
1233 CURLM *multi_handle = curl_multi_init();
1234 if (multi_handle == nullptr) {
1235 throw std::runtime_error("Failed to create curl multi-handle");
1236 }
1237
1238 int running_handles = 0;
1239 time_t last_maintenance = time(NULL);
1240 CURLMcode mres = CURLM_OK;
1241
1242 // Map from a file descriptor that has an outstanding broker request
1243 // to the corresponding CURL handle.
1244 std::unordered_map<int, WaitingForBroker> broker_reqs;
1245 std::vector<struct curl_waitfd> waitfds;
1246
1247 bool want_shutdown = false;
1248 while (!want_shutdown) {
1249 m_last_completed_cycle.store(std::chrono::system_clock::now().time_since_epoch().count());
1250 auto oldest_op = std::chrono::system_clock::now();
1251 for (const auto &entry : m_op_map) {
1252 OpRecord(*entry.second.first, OpKind::Update);
1253 if (entry.second.second < oldest_op) {
1254 oldest_op = entry.second.second;
1255 }
1256 }
1257 m_oldest_op.store(oldest_op.time_since_epoch().count());
1258
1259 // Try continuing any available handles that have more data
1260 while (true) {
1261 auto op = m_continue_queue->TryConsume();
1262 if (!op) {
1263 break;
1264 }
1265 // Avoid race condition where external thread added a continue operation to queue
1266 // while the curl worker thread failed the transfer.
1267 if (op->IsDone()) {
1268 m_logger->Debug(kLogXrdClHttp, "Ignoring continuation of operation that has already completed");
1269 continue;
1270 }
1271 m_logger->Debug(kLogXrdClHttp, "Continuing the curl handle from op %p on thread %d", op.get(), getthreadid());
1272 auto curl = op->GetCurlHandle();
1273 if (!op->ContinueHandle()) {
1274 op->Fail(XrdCl::errInternal, 0, "Failed to continue the curl handle for the operation");
1275 OpRecord(*op, OpKind::Error);
1276 op->ReleaseHandle();
1277 if (curl) {
1278 curl_multi_remove_handle(multi_handle, curl);
1279 curl_easy_cleanup(curl);
1280 m_op_map.erase(curl);
1281 }
1282 running_handles -= 1;
1283 continue;
1284 } else {
1285 auto iter = m_op_map.find(curl);
1286 if (iter != m_op_map.end()) iter->second.second = std::chrono::system_clock::now();
1287 }
1288 }
1289 // Consume from the shared new operation queue
1290 while (running_handles < static_cast<int>(m_max_ops)) {
1291 auto op = running_handles == 0 ? queue.Consume(std::chrono::seconds(1)) : queue.TryConsume();
1292 if (!op) {
1293 break;
1294 }
1295 auto curl = queue.GetHandle();
1296 if (curl == nullptr) {
1297 m_logger->Debug(kLogXrdClHttp, "Unable to allocate a curl handle");
1298 op->Fail(XrdCl::errInternal, ENOMEM, "Unable to get allocate a curl handle");
1299 continue;
1300 }
1301 try {
1302 auto rv = op->Setup(curl, *this);
1303 if (!rv) {
1304 m_logger->Debug(kLogXrdClHttp, "Failed to setup the curl handle");
1305 op->Fail(XrdCl::errInternal, ENOMEM, "Failed to setup the curl handle for the operation");
1306 continue;
1307 }
1308 if (!op->FinishSetup(curl)) {
1309 m_logger->Debug(kLogXrdClHttp, "Failed to finish setup of the curl handle");
1310 op->Fail(XrdCl::errInternal, ENOMEM, "Failed to finish setup of the curl handle for the operation");
1311 continue;
1312 }
1313 } catch (...) {
1314 m_logger->Debug(kLogXrdClHttp, "Unable to setup the curl handle");
1315 op->Fail(XrdCl::errInternal, ENOMEM, "Failed to setup the curl handle for the operation");
1316 continue;
1317 }
1318 op->SetContinueQueue(m_continue_queue);
1319
1320 if (op->IsDone()) {
1321 op->ReleaseHandle();
1322 queue.RecycleHandle(curl);
1323 continue;
1324 }
1325 m_op_map[curl] = {op, std::chrono::system_clock::now()};
1326
1327 // If the operation requires the result of the OPTIONS verb to function, then
1328 // we add that to the multi-handle instead, chaining the two calls together.
1329 if (op->RequiresOptions()) {
1330 std::string modified_url;
1331 std::shared_ptr<CurlOptionsOp> options_op(
1332 new CurlOptionsOp(
1333 curl, op,
1334 std::string(
1335 VerbsCache::GetUrlKey(op->GetUrl(), modified_url)
1336 ),
1337 m_logger, op->GetConnCalloutFunc()
1338 )
1339 );
1340 // Note this `curl` variable is not local to the conditional; it is the curl handle of the
1341 // CurlOptionsOp and will be added below to the multi-handle, causing it - not the parent's
1342 // curl handle - to be executed.
1343 curl = queue.GetHandle();
1344 if (curl == nullptr) {
1345 m_logger->Debug(kLogXrdClHttp, "Unable to allocate a curl handle");
1346 op->Fail(XrdCl::errInternal, ENOMEM, "Unable to get allocate a curl handle");
1347 OpRecord(*op, OpKind::Error);
1348 continue;
1349 }
1350 auto rv = options_op->Setup(curl, *this);
1351 if (!rv) {
1352 m_logger->Debug(kLogXrdClHttp, "Failed to allocate a curl handle for OPTIONS");
1353 continue;
1354 }
1355 m_op_map[curl] = {options_op, std::chrono::system_clock::now()};
1356 OpRecord(*options_op, OpKind::Start);
1357 running_handles += 1;
1358 } else {
1359 OpRecord(*op, OpKind::Start);
1360 }
1361
1362 auto mres = curl_multi_add_handle(multi_handle, curl);
1363 if (mres != CURLM_OK) {
1364 m_logger->Debug(kLogXrdClHttp, "Unable to add operation to the curl multi-handle");
1365 op->Fail(XrdCl::errInternal, mres, "Unable to add operation to the curl multi-handle");
1366 OpRecord(*op, OpKind::Error);
1367 continue;
1368 }
1369 m_logger->Debug(kLogXrdClHttp, "Added request for URL %s to worker thread for processing", op->GetUrl().c_str());
1370 running_handles += 1;
1371 }
1372
1373 // Maintain the periodic reporting of thread activity and fail any operations
1374 // that have expired / timed out.
1375 time_t now = time(NULL);
1376 time_t next_maintenance = last_maintenance + m_maintenance_period.load(std::memory_order_relaxed);
1377 if (now >= next_maintenance) {
1378 m_queue->Expire();
1379 m_continue_queue->Expire();
1380 m_logger->Debug(kLogXrdClHttp, "Curl worker thread %d is running %d operations",
1381 getthreadid(), running_handles);
1382 last_maintenance = now;
1383
1384 // Timeout all the pending broker requests.
1385 std::vector<std::pair<int, CURL *>> expired_ops;
1386 for (const auto &entry : broker_reqs) {
1387 if (entry.second.expiry < now) {
1388 expired_ops.emplace_back(entry.first, entry.second.curl);
1389 }
1390 }
1391 for (const auto &entry : expired_ops) {
1392 auto iter = m_op_map.find(entry.second);
1393 if (iter == m_op_map.end()) {
1394 m_logger->Warning(kLogXrdClHttp, "Found an expired curl handle with no corresponding operation!");
1395 } else {
1396
1397 CurlOptionsOp *options_op = nullptr;
1398 if ((options_op = dynamic_cast<CurlOptionsOp*>(iter->second.first.get())) != nullptr) {
1399 auto parent_op = options_op->GetOperation();
1400 bool parent_op_failed = false;
1401 if (parent_op->IsRedirect()) {
1402 std::string target;
1403 if (parent_op->Redirect(target) == CurlOperation::RedirectAction::Fail) {
1404 auto iter = m_op_map.find(options_op->GetParentCurlHandle());
1405 if (iter != m_op_map.end()) {
1406 OpRecord(*iter->second.first, OpKind::Error);
1407 iter->second.first->Fail(XrdCl::errErrorResponse, 0, "Failed to send OPTIONS to redirect target");
1408 m_op_map.erase(iter);
1409 running_handles -= 1;
1410 }
1411 parent_op_failed = true;
1412 } else {
1413 OpRecord(*parent_op, OpKind::Start);
1414 }
1415 } else {
1416 OpRecord(*parent_op, OpKind::Start);
1417 }
1418 if (!parent_op_failed){
1419 curl_multi_add_handle(multi_handle, options_op->GetParentCurlHandle());
1420 }
1421 }
1422
1423 iter->second.first->Fail(XrdCl::errConnectionError, 1, "Timeout: connection never provided for request");
1424 iter->second.first->ReleaseHandle();
1425 OpRecord(*(iter->second.first), OpKind::ConncallTimeout);
1426 m_op_map.erase(entry.second);
1427 curl_easy_cleanup(entry.second);
1428 running_handles -= 1;
1429 }
1430 broker_reqs.erase(entry.first);
1431 m_conncall_timeout.fetch_add(1, std::memory_order_relaxed);
1432 }
1433
1434 // Cleanup the fake connection cache entries.
1436 }
1437
1438 waitfds.clear();
1439 waitfds.resize(3 + broker_reqs.size());
1440
1441 waitfds[0].fd = queue.PollFD();
1442 waitfds[0].events = CURL_WAIT_POLLIN;
1443 waitfds[0].revents = 0;
1444 waitfds[1].fd = m_continue_queue->PollFD();
1445 waitfds[1].events = CURL_WAIT_POLLIN;
1446 waitfds[1].revents = 0;
1447 waitfds[2].fd = m_shutdown_pipe_r;
1448 waitfds[2].revents = 0;
1449 waitfds[2].events = CURL_WAIT_POLLIN | CURL_WAIT_POLLPRI;
1450
1451 int idx = 3;
1452 for (const auto &entry : broker_reqs) {
1453 waitfds[idx].fd = entry.first;
1454 waitfds[idx].events = CURL_WAIT_POLLIN|CURL_WAIT_POLLPRI;
1455 waitfds[idx].revents = 0;
1456 idx += 1;
1457 }
1458
1459 long timeo;
1460 curl_multi_timeout(multi_handle, &timeo);
1461 // These commented-out lines are purposely left; will need to revisit after the 0.9.1 release;
1462 // for now, they are too verbose on RHEL7.
1463 //m_logger->Debug(kLogXrdClHttp, "Curl advises a timeout of %ld ms", timeo);
1464 if (running_handles && timeo == -1) {
1465 // Bug workaround: we've seen RHEL7 libcurl have a race condition where it'll not
1466 // set a timeout while doing the DNS lookup; assume that if there are running handles
1467 // but no timeout, we've hit this bug.
1468 //m_logger->Debug(kLogXrdClHttp, "Will sleep for up to 50ms");
1469 mres = curl_multi_wait(multi_handle, &waitfds[0], waitfds.size(), 50, nullptr);
1470 } else {
1471 //m_logger->Debug(kLogXrdClHttp, "Will sleep for up to %d seconds", max_sleep_time);
1472 //mres = curl_multi_wait(multi_handle, &waitfds[0], waitfds.size(), max_sleep_time*1000, nullptr);
1473 // Temporary test: we've been seeing DNS lookups timeout on additional platforms. Switch to always
1474 // poll as curl_multi_wait doesn't seem to get notified when DNS lookups are done.
1475 mres = curl_multi_wait(multi_handle, &waitfds[0], waitfds.size(), 50, nullptr);
1476 }
1477 if (mres != CURLM_OK) {
1478 m_logger->Warning(kLogXrdClHttp, "Failed to wait on multi-handle: %d", mres);
1479 }
1480
1481 // Iterate through the waiting broker callbacks.
1482 for (const auto &entry : waitfds) {
1483 // Ignore the queue's poll fd.
1484 if (waitfds[0].fd == entry.fd || waitfds[1].fd == entry.fd) {
1485 continue;
1486 }
1487 // Handle shutdown requests
1488 if ((waitfds[2].fd == entry.fd) && entry.revents) {
1489 want_shutdown = true;
1490 break;
1491 }
1492 if ((entry.revents & CURL_WAIT_POLLIN) != CURL_WAIT_POLLIN) {
1493 continue;
1494 }
1495 auto handle = broker_reqs[entry.fd].curl;
1496 auto iter = m_op_map.find(handle);
1497 if (iter == m_op_map.end()) {
1498 m_logger->Warning(kLogXrdClHttp, "Internal error: broker responded on FD %d but no corresponding curl operation", entry.fd);
1499 broker_reqs.erase(entry.fd);
1500 m_conncall_errors.fetch_add(1, std::memory_order_relaxed);
1501 continue;
1502 }
1503 std::string err;
1504 auto result = iter->second.first->WaitSocketCallback(err);
1505 if (result == -1) {
1506 m_logger->Warning(kLogXrdClHttp, "Error when invoking the broker callback: %s", err.c_str());
1507
1508 CurlOptionsOp *options_op = nullptr;
1509 if ((options_op = dynamic_cast<CurlOptionsOp*>(iter->second.first.get())) != nullptr) {
1510 auto parent_op = options_op->GetOperation();
1511 bool parent_op_failed = false;
1512 if (parent_op->IsRedirect()) {
1513 std::string target;
1514 if (parent_op->Redirect(target) == CurlOperation::RedirectAction::Fail) {
1515 auto iter = m_op_map.find(options_op->GetParentCurlHandle());
1516 if (iter != m_op_map.end()) {
1517 OpRecord(*iter->second.first, OpKind::Error);
1518 iter->second.first->Fail(XrdCl::errErrorResponse, 0, "Failed to send OPTIONS to redirect target");
1519 m_op_map.erase(iter);
1520 running_handles -= 1;
1521 }
1522 parent_op_failed = true;
1523 } else {
1524 OpRecord(*parent_op, OpKind::Start);
1525 }
1526 } else {
1527 OpRecord(*parent_op, OpKind::Start);
1528 }
1529 if (!parent_op_failed){
1530 curl_multi_add_handle(multi_handle, options_op->GetParentCurlHandle());
1531 }
1532 }
1533
1534 iter->second.first->Fail(XrdCl::errErrorResponse, 1, err);
1535 OpRecord(*iter->second.first, OpKind::Error);
1536 m_op_map.erase(handle);
1537 broker_reqs.erase(entry.fd);
1538 m_conncall_errors.fetch_add(1, std::memory_order_relaxed);
1539 running_handles -= 1;
1540 } else {
1541 broker_reqs.erase(entry.fd);
1542 curl_multi_add_handle(multi_handle, handle);
1543 m_conncall_success.fetch_add(1, std::memory_order_relaxed);
1544 }
1545 }
1546
1547 // Do maintenance on the multi-handle
1548 int still_running;
1549 auto mres = curl_multi_perform(multi_handle, &still_running);
1550 if (mres == CURLM_CALL_MULTI_PERFORM) {
1551 continue;
1552 } else if (mres != CURLM_OK) {
1553 m_logger->Warning(kLogXrdClHttp, "Failed to perform multi-handle operation: %d", mres);
1554 break;
1555 }
1556
1557 CURLMsg *msg;
1558 do {
1559 int msgq = 0;
1560 msg = curl_multi_info_read(multi_handle, &msgq);
1561 if (msg && (msg->msg == CURLMSG_DONE)) {
1562 if (!msg->easy_handle) {
1563 m_logger->Warning(kLogXrdClHttp, "Logic error: got a callback for a null handle");
1564 mres = CURLM_BAD_EASY_HANDLE;
1565 break;
1566 }
1567 auto iter = m_op_map.find(msg->easy_handle);
1568 if (iter == m_op_map.end()) {
1569 m_logger->Error(kLogXrdClHttp, "Logic error: got a callback for an entry that doesn't exist");
1570 mres = CURLM_BAD_EASY_HANDLE;
1571 break;
1572 }
1573 auto op = iter->second.first;
1574 auto res = msg->data.result;
1575 bool keep_handle = false;
1576 bool waiting_on_callout = false;
1577 if (res == CURLE_OK) {
1578 auto sc = op->GetStatusCode();
1579 OpRecord(*op, OpKind::Finish);
1580 if (HTTPStatusIsError(sc)) {
1581 auto httpErr = HTTPStatusConvert(sc);
1582 op->Fail(httpErr.first, httpErr.second, op->GetStatusMessage());
1583 op->ReleaseHandle();
1584 // If this was a failed CurlOptionsOp, then we re-activate the parent handle.
1585 // If the parent handle was stopped at a redirect that now returns failure, then
1586 // we'll clean it up.
1587 CurlOptionsOp *options_op = nullptr;
1588 if ((options_op = dynamic_cast<CurlOptionsOp*>(op.get())) != nullptr) {
1589 auto parent_op = options_op->GetOperation();
1590 bool parent_op_failed = false;
1591 if (parent_op->IsRedirect()) {
1592 std::string target;
1593 if (parent_op->Redirect(target) == CurlOperation::RedirectAction::Fail) {
1594 OpRecord(*parent_op, OpKind::Error);
1595 m_op_map.erase(options_op->GetParentCurlHandle());
1596 running_handles -= 1;
1597 parent_op_failed = true;
1598 } else {
1599 OpRecord(*parent_op, OpKind::Start);
1600 }
1601 } else {
1602 OpRecord(*parent_op, OpKind::Start);
1603 }
1604 // Have curl execute the parent operation
1605 if (!parent_op_failed) {
1606 curl_multi_add_handle(multi_handle, options_op->GetParentCurlHandle());
1607 }
1608 }
1609 // The curl operation was successful, it's just the HTTP request failed; recycle the handle.
1610 queue.RecycleHandle(iter->first);
1611 } else {
1612 CurlOptionsOp *options_op = nullptr;
1613 // If this was a successful OPTIONS op, invoke the parent operation.
1614 if ((options_op = dynamic_cast<CurlOptionsOp*>(op.get()))) {
1615 options_op->Success();
1616 options_op->ReleaseHandle();
1617 // Note: op is scoped external to the conditional block
1618 op = options_op->GetOperation();
1619 op->OptionsDone();
1620 OpRecord(*op, OpKind::Start);
1621 curl_multi_add_handle(multi_handle, options_op->GetParentCurlHandle());
1622 curl_multi_remove_handle(multi_handle, iter->first);
1623 queue.RecycleHandle(iter->first);
1624 }
1625 // Check to see if the operation ended in a redirect (note: this might)
1626 // be invoked a second time if this was the parent operation of an OPTIONS
1627 // op.
1628 if (op->IsRedirect()) {
1629 std::string target;
1630 switch (op->Redirect(target)) {
1632 if (options_op) {
1633 // In this case, we failed immediately after an OPTIONS finished.
1634 // Since there's a Start recorded after the OPTIONS processing, we
1635 // must record an error.
1636 // In the non-OPTIONS case, we never recorded a second start and
1637 // don't need a matching failure.
1638 OpRecord(*op, OpKind::Error);
1639 }
1640 keep_handle = false;
1641 break;
1643 if (!options_op) {
1644 // In this case, the redirect occurred without any prior
1645 // OPTIONS call. This implies that `op` is the original call
1646 // and we need to restart it later and record another op start.
1647 keep_handle = true;
1648 OpRecord(*op, OpKind::Start);
1649 }
1650 break;
1652 {
1653 // The redirect resulted in a new endpoint where the cache lookup failed;
1654 // we need to know what HTTP verbs are in the server's Allow list before this
1655 // operation can continue. Inject a new CurlOptionsOp and chain it to the one
1656 // being processed. Once the OPTIONS request is done, then we'll restart this
1657 // operation.
1658 std::string modified_url;
1659 target = VerbsCache::GetUrlKey(target, modified_url);
1660 options_op = new CurlOptionsOp(iter->first, op, target, m_logger, op->GetConnCalloutFunc());
1661 std::shared_ptr<CurlOperation> new_op(options_op);
1662 auto curl = queue.GetHandle();
1663 if (curl == nullptr) {
1664 m_logger->Debug(kLogXrdClHttp, "Unable to allocate a curl handle");
1665 op->Fail(XrdCl::errInternal, ENOMEM, "Unable to get allocate a curl handle");
1666 keep_handle = false;
1667 options_op = nullptr;
1668 break;
1669 }
1670 OpRecord(*new_op, OpKind::Start);
1671 try {
1672 auto rv = new_op->Setup(curl, *this);
1673 if (!rv) {
1674 m_logger->Debug(kLogXrdClHttp, "Unable to configure a curl handle for OPTIONS");
1675 keep_handle = false;
1676 options_op = nullptr;
1677 break;
1678 }
1679 } catch (...) {
1680 m_logger->Debug(kLogXrdClHttp, "Unable to setup the curl handle for the OPTIONS operation");
1681 new_op->Fail(XrdCl::errInternal, ENOMEM, "Failed to setup the curl handle for the OPTIONS operation");
1682 OpRecord(*new_op, OpKind::Error);
1683 keep_handle = false;
1684 break;
1685 }
1686 new_op->SetContinueQueue(m_continue_queue);
1687 m_op_map[curl] = {new_op, std::chrono::system_clock::now()};
1688 auto mres = curl_multi_add_handle(multi_handle, curl);
1689 if (mres != CURLM_OK) {
1690 m_logger->Debug(kLogXrdClHttp, "Unable to add OPTIONS operation to the curl multi-handle: %s", curl_multi_strerror(mres));
1691 op->Fail(XrdCl::errInternal, mres, "Unable to add OPTIONS operation to the curl multi-handle");
1692 OpRecord(*new_op, OpKind::Error);
1693 break;
1694 }
1695 running_handles += 1;
1696 m_logger->Debug(kLogXrdClHttp, "Invoking the OPTIONS operation before redirect to %s", target.c_str());
1697 // The original curl operation needs to be kept around. Note that because options_op
1698 // is non-nil, we won't re-add the handle to the multi-handle.
1699 keep_handle = true;
1700 }
1701 }
1702 int callout_socket = op->WaitSocket();
1703 if ((waiting_on_callout = callout_socket >= 0)) {
1704 auto expiry = time(nullptr) + 20;
1705 m_logger->Debug(kLogXrdClHttp, "Creating a callout wait request on socket %d", callout_socket);
1706 broker_reqs[callout_socket] = {iter->first, expiry};
1707 m_conncall_req.fetch_add(1, std::memory_order_relaxed);
1708 }
1709 } else if (options_op) {
1710 // In this case, the OPTIONS call happened before the parent operation was started.
1711 curl_multi_add_handle(multi_handle, options_op->GetParentCurlHandle());
1712 }
1713 if (keep_handle) {
1714 curl_multi_remove_handle(multi_handle, iter->first);
1715 if (!waiting_on_callout && !options_op) {
1716 curl_multi_add_handle(multi_handle, iter->first);
1717 }
1718 } else if (!options_op) {
1719 // A multi-step operation may reset and reconfigure its
1720 // easy handle from Success(). Remove the completed
1721 // request before invoking it so libcurl no longer owns
1722 // the request configuration being replaced.
1723 curl_multi_remove_handle(multi_handle, iter->first);
1724 op->Success();
1725 if (op->IsDone()) {
1726 op->ReleaseHandle();
1727 // If the handle was successful, then we can recycle it.
1728 queue.RecycleHandle(iter->first);
1729 } else {
1730 // Multi-step operations may configure their easy handle
1731 // for another request from Success(). Keep ownership of
1732 // the handle and run the next request through this worker's
1733 // multi-handle like any other operation.
1734 auto next_res = curl_multi_add_handle(multi_handle, iter->first);
1735 if (next_res == CURLM_OK) {
1736 keep_handle = true;
1737 OpRecord(*op, OpKind::Start);
1738 } else {
1739 op->Fail(XrdCl::errInternal, next_res,
1740 "Unable to add the next operation request to the curl multi-handle");
1741 OpRecord(*op, OpKind::Error);
1742 op->ReleaseHandle();
1743 queue.RecycleHandle(iter->first);
1744 }
1745 }
1746 }
1747 }
1748 } else if (res == CURLE_COULDNT_CONNECT && op->UseConnectionCallout() && !op->GetTriedBoker()) {
1749 // In this case, we need to use the broker and the curl handle couldn't reuse
1750 // an existing socket.
1751 keep_handle = true;
1752 op->SetTriedBoker(); // Flag to ensure we try a connection only once per operation.
1753 std::string err;
1754 int wait_socket = -1;
1755 if (!op->StartConnectionCallout(err) || (wait_socket=op->WaitSocket()) == -1) {
1756 m_logger->Error(kLogXrdClHttp, "Failed to start broker-based connection: %s", err.c_str());
1757 op->ReleaseHandle();
1758 keep_handle = false;
1759 } else {
1760 curl_multi_remove_handle(multi_handle, iter->first);
1761 auto expiry = time(nullptr) + 20;
1762 m_logger->Debug(kLogXrdClHttp, "Curl operation requires a new TCP socket; waiting for callout to respond on socket %d", wait_socket);
1763 broker_reqs[wait_socket] = {iter->first, expiry};
1764 m_conncall_req.fetch_add(1, std::memory_order_relaxed);
1765 }
1766 } else {
1767 if (res == CURLE_ABORTED_BY_CALLBACK || res == CURLE_WRITE_ERROR) {
1768 // We cannot invoke the failure from within a callback as the curl thread and
1769 // original thread of execution may fight over the ownership of the handle memory.
1770 switch (op->GetError()) {
1772#ifdef HAVE_XPROTOCOL_TIMEREXPIRED
1773 op->Fail(XrdCl::errOperationExpired, 0, "Origin did not respond with headers within timeout");
1774#else
1775 op->Fail(XrdCl::errOperationExpired, 0, "Origin did not respond within timeout");
1776#endif
1777 OpRecord(*op, OpKind::Error);
1778 break;
1780 auto [ecode, emsg] = op->GetCallbackError();
1781 op->Fail(XrdCl::errErrorResponse, ecode, emsg);
1782 OpRecord(*op, OpKind::Error);
1783 break;
1784 }
1786 op->Fail(XrdCl::errOperationExpired, 0, "Operation timed out");
1787 OpRecord(*op, op->IsPaused() ? OpKind::ClientTimeout : OpKind::ServerTimeout);
1788 break;
1790 op->Fail(XrdCl::errOperationExpired, 0, "Transfer speed below minimum threshold");
1791 OpRecord(*op, OpKind::ServerTimeout);
1792 break;
1794 op->Fail(XrdCl::errOperationExpired, 0, "Transfer stalled for too long");
1795 OpRecord(*op, OpKind::ClientTimeout);
1796 break;
1798 op->Fail(XrdCl::errOperationExpired, 0, "Transfer stalled for too long");
1799 OpRecord(*op, OpKind::ServerTimeout);
1800 break;
1802 op->Fail(XrdCl::errInternal, 0, "Operation was aborted without recording an abort reason");
1803 OpRecord(*op, OpKind::Error);
1804 break;
1805 };
1806 CurlOptionsOp *options_op = nullptr;
1807 if ((options_op = dynamic_cast<CurlOptionsOp*>(op.get())) != nullptr) {
1808 auto parent_op = options_op->GetOperation();
1809 bool parent_op_failed = false;
1810 if (parent_op->IsRedirect()) {
1811 std::string target;
1812 if (parent_op->Redirect(target) == CurlOperation::RedirectAction::Fail) {
1813 auto iter = m_op_map.find(options_op->GetParentCurlHandle());
1814 if (iter != m_op_map.end()) {
1815 OpRecord(*iter->second.first, OpKind::Error);
1816 iter->second.first->Fail(XrdCl::errErrorResponse, 0, "Failed to send OPTIONS to redirect target");
1817 m_op_map.erase(iter);
1818 running_handles -= 1;
1819 }
1820 parent_op_failed = true;
1821 } else {
1822 OpRecord(*parent_op, OpKind::Start);
1823 }
1824 } else {
1825 OpRecord(*parent_op, OpKind::Start);
1826 }
1827 if (!parent_op_failed){
1828 curl_multi_add_handle(multi_handle, options_op->GetParentCurlHandle());
1829 }
1830 }
1831 } else {
1832 auto xrdCode = CurlCodeConvert(res);
1833 const auto curl_err = op->GetCurlErrorMessage();
1834 const char *curl_easy_err = curl_easy_strerror(res);
1835 const std::string fail_err = !curl_err.empty() ? curl_err : curl_easy_err;
1836 m_logger->Debug(kLogXrdClHttp, "Curl generated an error: %s (%d)", fail_err.c_str(), res);
1837 op->Fail(xrdCode.first, xrdCode.second, fail_err);
1838 OpRecord(*op, OpKind::Error);
1839 CurlOptionsOp *options_op = nullptr;
1840 if ((options_op = dynamic_cast<CurlOptionsOp*>(op.get())) != nullptr) {
1841 auto parent_op = options_op->GetOperation();
1842 bool parent_op_failed = false;
1843 if (parent_op->IsRedirect()) {
1844 std::string target;
1845 if (parent_op->Redirect(target) == CurlOperation::RedirectAction::Fail) {
1846 auto iter = m_op_map.find(options_op->GetParentCurlHandle());
1847 if (iter != m_op_map.end()) {
1848 OpRecord(*iter->second.first, OpKind::Error);
1849 iter->second.first->Fail(XrdCl::errErrorResponse, 0, "Failed to send OPTIONS to redirect target");
1850 m_op_map.erase(iter);
1851 running_handles -= 1;
1852 }
1853 parent_op_failed = true;
1854 }
1855 }
1856 if (!parent_op_failed){
1857 curl_multi_add_handle(multi_handle, options_op->GetParentCurlHandle());
1858 }
1859 }
1860 }
1861 op->ReleaseHandle();
1862 }
1863 if (!keep_handle) {
1864 curl_multi_remove_handle(multi_handle, iter->first);
1865 if (res != CURLE_OK) {
1866 curl_easy_cleanup(iter->first);
1867 }
1868 for (auto &req : broker_reqs) {
1869 if (req.second.curl == iter->first) {
1870 m_logger->Warning(kLogXrdClHttp, "Curl handle finished while a broker operation was outstanding");
1871 m_conncall_errors.fetch_add(1, std::memory_order_relaxed);
1872 }
1873 }
1874 m_op_map.erase(iter);
1875 running_handles -= 1;
1876 }
1877 }
1878 } while (msg);
1879 }
1880
1881 for (auto map_entry : m_op_map) {
1882 if (mres) {
1883 map_entry.second.first->Fail(XrdCl::errInternal, mres, curl_multi_strerror(mres));
1884 OpRecord(*map_entry.second.first, OpKind::Error);
1885 }
1886 if (multi_handle && map_entry.first) curl_multi_remove_handle(multi_handle, map_entry.first);
1887 }
1888
1889 m_queue->ReleaseHandles();
1890 curl_multi_cleanup(multi_handle);
1891}
1892
1893void
1894CurlWorker::Shutdown()
1895{
1896 m_queue->Shutdown();
1897 if (m_shutdown_pipe_w == -1) {
1898 m_logger->Debug(kLogXrdClHttp, "Curl worker shutdown prior to launch of thread");
1899 return;
1900 }
1901 close(m_shutdown_pipe_w);
1902 m_shutdown_pipe_w = -1;
1903
1904 // wait for worker thread to exit
1905 m_self_tid.join();
1906
1907 {
1908 std::unique_lock lk(m_worker_stats_mutex);
1909 m_workers_last_completed_cycle[m_stats_offset] = nullptr;
1910 m_workers_oldest_op[m_stats_offset] = nullptr;
1911 }
1912 m_logger->Debug(kLogXrdClHttp, "Curl worker thread shutdown has completed.");
1913}
1914
1915void
1916CurlWorker::ShutdownAll()
1917{
1918 std::unique_lock lock(m_workers_mutex);
1919 for (auto &worker : m_workers) {
1920 worker->Shutdown();
1921 }
1922}
1923
1924CurlWorker::initcontrol::initcontrol()
1925{
1926 curl_global_init(CURL_GLOBAL_DEFAULT);
1927}
1928
1929CurlWorker::initcontrol::~initcontrol()
1930{
1931 ShutdownAll();
1932 curl_global_cleanup();
1933}
@ kXR_InvalidRequest
@ kXR_Impossible
@ kXR_TimerExpired
@ kXR_NotAuthorized
@ kXR_NotFound
@ kXR_FileLocked
@ kXR_overQuota
@ kXR_Unsupported
@ kXR_Conflict
@ kXR_ServerError
@ kXR_Overloaded
@ kXR_ReqTimedOut
std::pair< uint16_t, uint32_t > CurlCodeConvert(CURLcode res)
void CURL
std::string obfuscateAuth(const std::string &input)
#define close(a)
Definition XrdPosix.hh:48
#define write(a, b, c)
Definition XrdPosix.hh:121
#define read(a, b, c)
Definition XrdPosix.hh:86
int emsg(int rc, char *msg)
char Value[ValuSize]
Definition XrdCksData.hh:53
int Set(const char *csName)
Definition XrdCksData.hh:81
bool Set(ChecksumType ctype, const std::array< unsigned char, g_max_checksum_length > &value)
virtual void Success()=0
bool FinishSetup(CURL *curl)
const std::string & GetUrl() const
CURL * GetCurlHandle() const
static const std::string GetVerbString(HttpVerb)
virtual HttpVerb GetVerb() const =0
std::string GetCurlErrorMessage() const
virtual void ReleaseHandle()
virtual bool RequiresOptions() const
static void CleanupDnsCache()
std::tuple< uint64_t, std::chrono::steady_clock::duration, std::chrono::steady_clock::duration, std::chrono::steady_clock::duration > StatisticsReset()
virtual bool ContinueHandle()
std::string GetStatusMessage() const
CreateConnCalloutType GetConnCalloutFunc() const
XrdClHttp::HttpVerb HttpVerb
virtual void Fail(uint16_t errCode, uint32_t errNum, const std::string &)
virtual RedirectAction Redirect(std::string &target)
virtual void SetContinueQueue(std::shared_ptr< XrdClHttp::HandlerQueue > queue)
bool StartConnectionCallout(std::string &err)
virtual bool Setup(CURL *curl, CurlWorker &)
std::pair< XErrorCode, std::string > GetCallbackError() const
CurlOptionsOp(CURL *curl, std::shared_ptr< CurlOperation > op, const std::string &url, XrdCl::Log *log, CreateConnCalloutType callout)
std::shared_ptr< CurlOperation > GetOperation() const
CURL * GetParentCurlHandle() const
void Fail(uint16_t errCode, uint32_t errNum, const std::string &) override
std::tuple< std::string, std::string > ClientX509CertKeyFile() const
CurlWorker(std::shared_ptr< HandlerQueue > queue, VerbsCache &cache, XrdCl::Log *logger)
static void RunStatic(CurlWorker *myself)
void Start(std::unique_ptr< XrdClHttp::CurlWorker > self, std::thread tid)
static std::string GetMonitoringJson()
std::shared_ptr< CurlOperation > Consume(std::chrono::steady_clock::duration)
HandlerQueue(unsigned max_pending_ops)
void Produce(std::shared_ptr< CurlOperation > handler)
static std::string GetMonitoringJson()
std::shared_ptr< CurlOperation > TryConsume()
void SetMultipartSeparator(const std::string_view &sep)
static void ParseDigest(const std::string &digest, XrdClHttp::ChecksumInfo &info)
static bool Canonicalize(std::string &headerName)
bool Parse(const std::string &headers)
static std::string ChecksumTypeToDigestName(XrdClHttp::ChecksumType type)
static bool Base64Decode(std::string_view input, std::array< unsigned char, g_max_checksum_length > &output)
static std::string_view GetUrlKey(const std::string &url, std::string &modified_url)
static Log * GetLog()
Get default log.
static Env * GetEnv()
Get default client environment.
bool GetInt(const std::string &key, int &value)
Definition XrdClEnv.cc:115
Handle diagnostics.
Definition XrdClLog.hh:101
@ DumpMsg
print details of the request and responses
Definition XrdClLog.hh:113
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
URL representation.
Definition XrdClURL.hh:31
std::string GetURL() const
Get the URL.
Definition XrdClURL.hh:86
const ParamsMap & GetParams() const
Get the URL params.
Definition XrdClURL.hh:244
static void splitString(Container &result, const std::string &input, const std::string &delimiter)
Split a string.
Definition XrdClUtils.hh:56
static void Trim(std::string &str)
Trim a string.
std::pair< uint16_t, uint32_t > HTTPStatusConvert(unsigned status)
size_t GetChecksumLength(ChecksumType ctype)
void InjectBearerToken(const XrdCl::URL &url, std::vector< std::pair< std::string, std::string > > &headers, XrdCl::Log *logger=nullptr)
CURL * GetHandle(bool verbose)
bool HTTPStatusIsError(unsigned status)
bool ShouldUseBearerToken(const std::string &protocols, bool hasX509Credential, bool hasBearerToken)
std::string_view ltrim_view(const std::string_view &input_view)
const uint64_t kLogXrdClHttp
std::string GetBearerToken(XrdCl::Log *logger=nullptr)
void ConfigureHandle(CURL *curl, bool verbose)
std::string_view trim_view(const std::string_view &input_view)
const uint16_t errUnknown
Unknown error.
const uint16_t errInvalidAddr
const uint16_t errRedirectLimit
const uint16_t errErrorResponse
const uint16_t errTlsError
const uint16_t errOperationExpired
const uint16_t errLoginFailed
const uint16_t errDataError
data is corrupted
const uint16_t errInternal
Internal error.
const uint16_t errInvalidArgs
const uint16_t errConnectionError
const uint16_t errNotSupported
const uint16_t errSocketError
const uint16_t errCorruptedHeader
const uint16_t errNone
No error.