XRootD
XrdHttpTpcTPC.cc
Go to the documentation of this file.
2 #include "XrdNet/XrdNetAddr.hh"
3 #include "XrdNet/XrdNetUtils.hh"
4 #include "XrdOuc/XrdOucEnv.hh"
5 #include "XrdSec/XrdSecEntity.hh"
8 #include "XrdSys/XrdSysFD.hh"
9 #include "XrdVersion.hh"
10 
12 #include "XrdOuc/XrdOucTUtils.hh"
14 #include "XrdHttp/XrdHttpUtils.hh"
15 
16 #include <curl/curl.h>
17 
18 #include <dlfcn.h>
19 #include <fcntl.h>
20 
21 #include <algorithm>
22 #include <memory>
23 #include <sstream>
24 #include <stdexcept>
25 #include <thread>
26 
27 #include "XrdHttpTpcState.hh"
28 #include "XrdHttpTpcStream.hh"
29 #include "XrdHttpTpcTPC.hh"
30 #include <fstream>
31 
32 using namespace TPC;
33 
34 XrdXrootdTpcMon* TPCHandler::TPCLogRecord::tpcMonitor = 0;
35 
36 uint64_t TPCHandler::m_monid{0};
37 int TPCHandler::m_marker_period = 5;
38 size_t TPCHandler::m_block_size = 16*1024*1024;
39 size_t TPCHandler::m_small_block_size = 1*1024*1024;
40 XrdSysMutex TPCHandler::m_monid_mutex;
41 bool TPCHandler::allowMissingCRL = false;
42 
44 
45 /******************************************************************************/
46 /* T P C H a n d l e r : : T P C L o g R e c o r d D e s t r u c t o r */
47 /******************************************************************************/
48 
49 TPCHandler::TPCLogRecord::~TPCLogRecord()
50 {
51 // Record monitoring data is enabled
52 //
53  if (tpcMonitor)
54  {XrdXrootdTpcMon::TpcInfo monInfo;
55 
56  monInfo.clID = clID.c_str();
57  monInfo.begT = begT;
58  gettimeofday(&monInfo.endT, 0);
59 
60  if (mTpcType == TpcType::Pull)
61  {monInfo.dstURL = local.c_str();
62  monInfo.srcURL = remote.c_str();
63  } else {
64  monInfo.dstURL = remote.c_str();
65  monInfo.srcURL = local.c_str();
67  }
68 
69  if (!status) monInfo.endRC = 0;
70  else if (tpc_status > 0) monInfo.endRC = tpc_status;
71  else monInfo.endRC = 1;
72  monInfo.strm = static_cast<unsigned char>(streams);
73  monInfo.fSize = (bytes_transferred < 0 ? 0 : bytes_transferred);
74  if (!isIPv6) monInfo.opts |= XrdXrootdTpcMon::TpcInfo::isIPv4;
75 
76  tpcMonitor->Report(monInfo);
77  }
78 }
79 
80 /******************************************************************************/
81 /* C u r l D e l e t e r : : o p e r a t o r ( ) */
82 /******************************************************************************/
83 
85 {
86  if (curl) curl_easy_cleanup(curl);
87 }
88 
89 /******************************************************************************/
90 /* s o c k o p t _ s e t c l o e x e c _ c a l l b a c k */
91 /******************************************************************************/
92 
101 int TPCHandler::sockopt_callback(void *clientp, curl_socket_t curlfd, curlsocktype purpose) {
102  TPCLogRecord * rec = (TPCLogRecord *)clientp;
103  if (purpose == CURLSOCKTYPE_IPCXN && rec && rec->pmarkManager.isEnabled()) {
104  // We will not reach this callback if the corresponding socket could not have been connected
105  // the socket is already connected only if the packet marking is enabled
106  return CURL_SOCKOPT_ALREADY_CONNECTED;
107  }
108  return CURL_SOCKOPT_OK;
109 }
110 
111 /******************************************************************************/
112 /* o p e n s o c k e t _ c a l l b a c k */
113 /******************************************************************************/
114 
115 
120 int TPCHandler::opensocket_callback(void *clientp,
121  curlsocktype purpose,
122  struct curl_sockaddr *aInfo)
123 {
124  /* CURLSOCKTYPE_IPCXN (for IP based connections) is the only type currently known by curl,
125  * so let's make sure to reject other types if they appear in the furure */
126  if (purpose != CURLSOCKTYPE_IPCXN)
127  return CURL_SOCKET_BAD;
128 
129  if (!aInfo)
130  return CURL_SOCKET_BAD;
131 
132  // Create the socket (note that O_CLOEXEC flag will be set)
133  int fd = XrdSysFD_Socket(aInfo->family, aInfo->socktype, aInfo->protocol);
134 
135  if (fd < 0) {
136  return CURL_SOCKET_BAD;
137  }
138 
139  if (!clientp)
140  return fd;
141 
142  XrdNetAddr thePeer(&(aInfo->addr));
143  TPCLogRecord *rec = static_cast<TPCLogRecord*>(clientp);
144 
145  /* Reject attempts to connect to local/private addresses unless allowed by configuration */
146  if ((!rec->allow_private && thePeer.isPrivate()) || (!rec->allow_local && thePeer.isLocal())) {
147  rec->tpc_status = 403; // Forbidden
148  rec->m_log->Emsg(rec->log_prefix.c_str(),
149  "Connection to local/private address is forbidden");
150  close(fd);
151  return CURL_SOCKET_BAD;
152  }
153 
154  rec->isIPv6 = (thePeer.isIPType(XrdNetAddrInfo::IPv6) && !thePeer.isMapped());
155 
156  std::stringstream connectErrMsg;
157  if(!rec->pmarkManager.connect(fd, &(aInfo->addr), aInfo->addrlen, CONNECT_TIMEOUT, connectErrMsg)) {
158  // at this point fd has already been closed
159  rec->m_log->Emsg(rec->log_prefix.c_str(), "Unable to connect socket: ", connectErrMsg.str().c_str());
160  return CURL_SOCKET_BAD;
161  }
162 
163  return fd;
164 }
165 
166 int TPCHandler::closesocket_callback(void *clientp, curl_socket_t fd) {
167  TPCLogRecord * rec = (TPCLogRecord *)clientp;
168 
169  // Destroy the PMark handle associated to the file descriptor before closing it.
170  // Otherwise, we would lose the socket usage information if the socket is closed before
171  // the PMark handle is closed.
172  rec->pmarkManager.endPmark(fd);
173 
174  return close(fd);
175 }
176 
177 /******************************************************************************/
178 /* s s l _ c t x _ c a l l b a c k */
179 /******************************************************************************/
180 
187 int TPCHandler::ssl_ctx_callback(CURL *curl, void *ssl_ctx, void *clientp) {
188  TPCLogRecord * rec = (TPCLogRecord *)clientp;
189  SSL_CTX* ctx = static_cast<SSL_CTX*>(ssl_ctx);
190 
191  if (rec && rec->ca_store) {
192  // Bumps the store's reference count instead of re-parsing the CA and CRL
193  // bundles for this connection. libcurl runs this callback after it has
194  // applied its own TLS options, so this replaces whatever store it built.
195  SSL_CTX_set1_cert_store(ctx, rec->ca_store.get());
196  }
197  if (allowMissingCRL) {
198  // verify_callback only excuses X509_V_ERR_UNABLE_TO_GET_CRL, i.e. a CA in
199  // the chain for which no CRL could be found. Every other verification rule
200  // still applies, including revocation itself whenever a CRL is present.
201  SSL_CTX_set_verify(ctx, SSL_VERIFY_PEER, verify_callback);
202  }
203  return CURLE_OK;
204 }
205 
206 int TPCHandler::verify_callback(int preverify_ok, X509_STORE_CTX* ctx) {
207  if (preverify_ok == 1) return 1;
208 
209  int err = X509_STORE_CTX_get_error(ctx);
210 
211  if (err == X509_V_ERR_UNABLE_TO_GET_CRL) {
212  X509_STORE_CTX_set_error(ctx, X509_V_OK);
213  return 1;
214  }
215 
216  return 0;
217 }
218 
219 /******************************************************************************/
220 /* p r e p a r e U R L */
221 /******************************************************************************/
222 
223 // See XrdHttpTpcUtils::prepareOpenURL() documentation
224 std::string TPCHandler::prepareURL(XrdHttpExtReq &req) {
225  return XrdHttpTpcUtils::prepareOpenURL(req.resource, req.headers,hdr2cgimap);
226 }
227 
228 /******************************************************************************/
229 /* e n c o d e _ x r o o t d _ o p a q u e _ t o _ u r i */
230 /******************************************************************************/
231 
232 // When processing a redirection from the filesystem layer, it is permitted to return
233 // some xrootd opaque data. The quoting rules for xrootd opaque data are significantly
234 // more permissive than a URI (basically, only '&' and '=' are disallowed while some
235 // URI parsers may dislike characters like '"'). This function takes an opaque string
236 // (e.g., foo=1&bar=2&baz=") and makes it safe for all URI parsers.
237 std::string encode_xrootd_opaque_to_uri(CURL *curl, const std::string &opaque)
238 {
239  std::stringstream parser(opaque);
240  std::string sequence;
241  std::stringstream output;
242  bool first = true;
243  while (getline(parser, sequence, '&')) {
244  if (sequence.empty()) {continue;}
245  size_t equal_pos = sequence.find('=');
246  char *val = NULL;
247  if (equal_pos != std::string::npos)
248  val = curl_easy_escape(curl, sequence.c_str() + equal_pos + 1, sequence.size() - equal_pos - 1);
249  // Do not emit parameter if value exists and escaping failed.
250  if (!val && equal_pos != std::string::npos) {continue;}
251 
252  if (!first) output << "&";
253  first = false;
254  output << sequence.substr(0, equal_pos);
255  if (val) {
256  output << "=" << val;
257  curl_free(val);
258  }
259  }
260  return output.str();
261 }
262 
263 /******************************************************************************/
264 /* T P C H a n d l e r : : C o n f i g u r e C u r l C A */
265 /******************************************************************************/
266 
267 bool
268 TPCHandler::ConfigureCurlCA(CURL *curl, TPCLogRecord &rec)
269 {
270  // Preferred path: hand libcurl the CA/CRL store that XrdTlsTempCA already
271  // parsed, rather than the bundle filenames. Passing filenames makes libcurl
272  // build a private X509_STORE per connection, which costs tens of MB for a grid
273  // CA directory and is held for the whole transfer; sharing one store makes that
274  // a reference count. See https://github.com/xrootd/xrootd/issues/2873
275  //
276  // Skipped when m_cafile is set, so that the http.cafile precedence established
277  // at the bottom of this function is preserved.
278  if (m_ca_file && m_sslctx_supported && m_cafile.empty()) {
279  rec.ca_store = m_ca_file->CAStore();
280  if (!rec.ca_store) {
281  m_log.Log(Error, "TpcHandler", "No CA store is available; refusing to "
282  "fall back to libcurl's default CA bundle");
283  return false;
284  }
285  // Stop libcurl loading its build-time default bundle, which the callback
286  // below would only discard; the callback supplies the trust anchors.
287  curl_easy_setopt(curl, CURLOPT_CAINFO, static_cast<char *>(nullptr));
288  curl_easy_setopt(curl, CURLOPT_CAPATH, static_cast<char *>(nullptr));
289  curl_easy_setopt(curl, CURLOPT_SSL_CTX_FUNCTION, ssl_ctx_callback);
290  curl_easy_setopt(curl, CURLOPT_SSL_CTX_DATA, &rec);
291  return true;
292  }
293 
294  auto ca_filename = m_ca_file ? m_ca_file->CAFilename() : "";
295  auto crl_filename = m_ca_file ? m_ca_file->CRLFilename() : "";
296  if (!ca_filename.empty() && !crl_filename.empty()) {
297  curl_easy_setopt(curl, CURLOPT_CAINFO, ca_filename.c_str());
298  //Check that the CRL file contains at least one entry before setting this option to curl
299  //Indeed, an empty CRL file will make curl unhappy and therefore will fail
300  //all HTTP TPC transfers (https://github.com/xrootd/xrootd/issues/1543)
301  std::ifstream in(crl_filename, std::ifstream::ate | std::ifstream::binary);
302  if(in.tellg() > 0 && m_ca_file->atLeastOneValidCRLFound()){
303  curl_easy_setopt(curl, CURLOPT_CRLFILE, crl_filename.c_str());
304  if (allowMissingCRL) {
305  // No need to set the callback if there is no need to do it
306  curl_easy_setopt(curl, CURLOPT_SSL_CTX_FUNCTION, ssl_ctx_callback);
307  }
308  } else {
309  std::ostringstream oss;
310  oss << "No valid CRL file has been found in the file " << crl_filename << ". Disabling CRL checking.";
311  m_log.Log(Warning,"TpcHandler",oss.str().c_str());
312  }
313  }
314  else if (!m_cadir.empty()) {
315  curl_easy_setopt(curl, CURLOPT_CAPATH, m_cadir.c_str());
316  }
317  if (!m_cafile.empty()) {
318  curl_easy_setopt(curl, CURLOPT_CAINFO, m_cafile.c_str());
319  }
320  return true;
321 }
322 
323 void
324 TPCHandler::ConfigureCurlLowSpeed(CURL *curl)
325 {
326  // Older versions have poor transfer performance when low-speed limits are
327  // enabled; this was corrected in curl commit cacdc27f for version 7.38.0.
328  curl_version_info_data *curl_ver = curl_version_info(CURLVERSION_NOW);
329  if (m_low_speed_limit > 0 && curl_ver && curl_ver->age > 0 &&
330  curl_ver->version_num >= 0x072600) {
331  curl_easy_setopt(curl, CURLOPT_LOW_SPEED_TIME, m_low_speed_time);
332  curl_easy_setopt(curl, CURLOPT_LOW_SPEED_LIMIT, m_low_speed_limit);
333  }
334 }
335 
336 
337 bool TPCHandler::MatchesPath(const char *verb, const char *path) {
338  return !strcmp(verb, "COPY") || !strcmp(verb, "OPTIONS");
339 }
340 
341 /******************************************************************************/
342 /* P r e p a r e U R L */
343 /******************************************************************************/
344 
345 static std::string PrepareURL(const std::string &url)
346 {
347  const std::string replace_schemes[] = { "davs://", "s3://", "s3s://" };
348 
349  for (const auto& s : replace_schemes)
350  if (url.compare(0, s.size(), s) == 0)
351  return "https://" + url.substr(s.size());
352 
353  return url;
354 }
355 
356 static bool IsAllowedScheme(const std::string& url)
357 {
358  const std::string allowed_schemes[] = { "https://", "http://" };
359 
360  for (const auto& s : allowed_schemes)
361  if (url.compare(0, s.size(), s) == 0)
362  return true;
363 
364  return false;
365 }
366 
367 /******************************************************************************/
368 /* T P C H a n d l e r : : P r o c e s s R e q */
369 /******************************************************************************/
370 
372  if (req.verb == "OPTIONS") {
373  return ProcessOptionsReq(req);
374  }
375  auto header = XrdOucTUtils::caseInsensitiveFind(req.headers,"credential");
376  if (header != req.headers.end()) {
377  if (header->second != "none") {
378  m_log.Emsg("ProcessReq", "COPY requested an unsupported credential type: ", header->second.c_str());
379  return req.SendSimpleResp(400, NULL, NULL, "COPY requestd an unsupported Credential type", 0);
380  }
381  }
382  auto srcHeader = XrdOucTUtils::caseInsensitiveFind(req.headers,"source");
383  auto dstHeader = XrdOucTUtils::caseInsensitiveFind(req.headers,"destination");
384  // A Source header asks for a pull, a Destination header for a push; asking
385  // for both at once is ambiguous, so the request is rejected instead of
386  // arbitrarily honouring one of the two.
387  if (srcHeader != req.headers.end() && dstHeader != req.headers.end()) {
388  const char *error_both = "COPY rejected: both a Source and a Destination header were specified";
389  m_log.Emsg("ProcessReq", error_both);
390  return req.SendSimpleResp(400, NULL, NULL, error_both, 0);
391  }
392  if (srcHeader != req.headers.end()) {
393  std::string src = PrepareURL(srcHeader->second);
394  if (!IsAllowedScheme(src)) {
395  const char *error_src = "COPY rejected: disallowed scheme in source URL";
396  m_log.Emsg("ProcessReq", error_src, src.c_str());
397  return req.SendSimpleResp(400, NULL, NULL, error_src, 0);
398  }
399  return ProcessPullReq(src, req);
400  }
401  if (dstHeader != req.headers.end()) {
402  const std::string& dst = dstHeader->second;
403  if (!IsAllowedScheme(dst)) {
404  const char *error_dst = "COPY rejected: disallowed scheme in destination URL";
405  m_log.Emsg("ProcessReq", error_dst, dst.c_str());
406  return req.SendSimpleResp(400, NULL, NULL, error_dst, 0);
407  }
408  return ProcessPushReq(dst, req);
409  }
410  m_log.Emsg("ProcessReq", "COPY verb requested but no source or destination specified.");
411  return req.SendSimpleResp(400, NULL, NULL, "No Source or Destination specified", 0);
412 }
413 
414 /******************************************************************************/
415 /* T P C H a n d l e r D e s t r u c t o r */
416 /******************************************************************************/
417 
419  m_sfs = NULL;
420 }
421 
422 /******************************************************************************/
423 /* T P C H a n d l e r C o n s t r u c t o r */
424 /******************************************************************************/
425 
426 TPCHandler::TPCHandler(XrdSysError *log, const char *config, XrdOucEnv *myEnv) :
427  m_allow_local(false),
428  m_allow_private(true),
429  m_desthttps(false),
430  m_fixed_route(false),
431  m_low_speed_limit(10*1024),
432  m_low_speed_time(2*60),
433  m_timeout(60),
434  m_first_timeout(120),
435  m_log(log->logger(), "TPC_"),
436  m_sfs(NULL)
437 {
438  if (!Configure(config, myEnv)) {
439  throw std::runtime_error("Failed to configure the HTTP third-party-copy handler.");
440  }
441 
442 // Extract out the TPC monitoring object (we share it with xrootd).
443 //
444  XrdXrootdGStream *gs = (XrdXrootdGStream*)myEnv->GetPtr("Tpc.gStream*");
445  if (gs)
446  TPCLogRecord::tpcMonitor = new XrdXrootdTpcMon("http",log->logger(),*gs);
447 }
448 
449 /******************************************************************************/
450 /* T P C H a n d l e r : : P r o c e s s O p t i o n s R e q */
451 /******************************************************************************/
452 
456 int TPCHandler::ProcessOptionsReq(XrdHttpExtReq &req) {
457  return req.SendSimpleResp(200, NULL, (char *) "DAV: 1\r\nDAV: <http://apache.org/dav/propset/fs/1>\r\nAllow: HEAD,GET,PUT,PROPFIND,DELETE,OPTIONS,COPY", NULL, 0);
458 }
459 
460 /******************************************************************************/
461 /* T P C H a n d l e r : : G e t A u t h z */
462 /******************************************************************************/
463 
464 std::string TPCHandler::GetAuthz(XrdHttpExtReq &req) {
465  std::string authz;
466  auto authz_header = XrdOucTUtils::caseInsensitiveFind(req.headers,"authorization");
467  if (authz_header != req.headers.end()) {
468  std::stringstream ss;
469  ss << "authz=" << encode_str(authz_header->second);
470  authz += ss.str();
471  }
472  return authz;
473 }
474 
475 /******************************************************************************/
476 /* T P C H a n d l e r : : R e d i r e c t T r a n s f e r */
477 /******************************************************************************/
478 
479 int TPCHandler::RedirectTransfer(CURL *curl, const std::string &redirect_resource,
480  XrdHttpExtReq &req, XrdOucErrInfo &error, TPCLogRecord &rec)
481 {
482  int port;
483  const char *ptr = error.getErrText(port);
484  if ((ptr == NULL) || (*ptr == '\0') || (port == 0)) {
485  rec.status = 500;
486  std::stringstream ss;
487  ss << "Internal error: redirect without hostname";
488  logTransferEvent(LogMask::Error, rec, "REDIRECT_INTERNAL_ERROR", ss.str());
489  return req.SendSimpleResp(rec.status, NULL, NULL, generateClientErr(ss, rec).c_str(), 0);
490  }
491 
492  // Construct redirection URL taking into consideration any opaque info
493  std::string rdr_info = ptr;
494  std::string host, opaque;
495  size_t pos = rdr_info.find('?');
496  host = rdr_info.substr(0, pos);
497 
498  if (pos != std::string::npos) {
499  opaque = rdr_info.substr(pos + 1);
500  }
501 
502  std::stringstream ss;
503  ss << "Location: http" << (m_desthttps ? "s" : "") << "://" << host << ":" << port << "/" << redirect_resource;
504 
505  if (!opaque.empty()) {
506  ss << "?" << encode_xrootd_opaque_to_uri(curl, opaque);
507  }
508 
509  rec.status = 307;
510  logTransferEvent(LogMask::Info, rec, "REDIRECT", ss.str());
511  return req.SendSimpleResp(rec.status, NULL, const_cast<char *>(ss.str().c_str()),
512  NULL, 0);
513 }
514 
515 /******************************************************************************/
516 /* T P C H a n d l e r : : O p e n W a i t S t a l l */
517 /******************************************************************************/
518 
519 int TPCHandler::OpenWaitStall(XrdSfsFile &fh, const std::string &resource,
520  int mode, int openMode, const XrdSecEntity &sec,
521  const std::string &authz)
522 {
523  int open_result;
524  while (1) {
525  int orig_ucap = fh.error.getUCap();
526  fh.error.setUCap(orig_ucap | XrdOucEI::uIPv64);
527  std::string opaque;
528  size_t pos = resource.find('?');
529  // Extract the path and opaque info from the resource
530  std::string path = resource.substr(0, pos);
531 
532  if (pos != std::string::npos) {
533  opaque = resource.substr(pos + 1);
534  }
535 
536  // Append the authz information if there are some
537  if(!authz.empty()) {
538  opaque += (opaque.empty() ? "" : "&");
539  opaque += authz;
540  }
541  open_result = fh.open(path.c_str(), mode, openMode, &sec, opaque.c_str());
542 
543  if ((open_result == SFS_STALL) || (open_result == SFS_STARTED)) {
544  int secs_to_stall = fh.error.getErrInfo();
545  if (open_result == SFS_STARTED) {secs_to_stall = secs_to_stall/2 + 5;}
546  std::this_thread::sleep_for (std::chrono::seconds(secs_to_stall));
547  }
548  break;
549  }
550  return open_result;
551 }
552 
553 /******************************************************************************/
554 /* T P C H a n d l e r : : D e t e r m i n e X f e r S i z e */
555 /******************************************************************************/
556 
557 
558 
562 int TPCHandler::DetermineXferSize(CURL *curl, XrdHttpExtReq &req, State &state,
563  bool &success, TPCLogRecord &rec, bool shouldReturnErrorToClient) {
564  success = false;
565  curl_easy_setopt(curl, CURLOPT_NOBODY, 1);
566  // Set a custom timeout of 60 seconds (= CONNECT_TIMEOUT for convenience) for the HEAD request
567  curl_easy_setopt(curl, CURLOPT_TIMEOUT, CONNECT_TIMEOUT);
568  CURLcode res;
569  res = curl_easy_perform(curl);
570  //Immediately set the CURLOPT_NOBODY flag to 0 as we anyway
571  //don't want the next curl call to do be a HEAD request
572  curl_easy_setopt(curl, CURLOPT_NOBODY, 0);
573  // Reset the CURLOPT_TIMEOUT to no timeout (default)
574  curl_easy_setopt(curl, CURLOPT_TIMEOUT, 0L);
575  curl_easy_setopt(curl, CURLOPT_FAILONERROR, true);
576 
577  std::stringstream ss;
578 
579  if (state.GetStatusCode() >= 400)
580  res = CURLE_HTTP_RETURNED_ERROR;
581 
582  if (res != CURLE_OK) { /* curl failed */
583  ss << curl_easy_strerror(res);
584  switch (res) {
585  case CURLE_HTTP_RETURNED_ERROR: /* remote side may have returned an error */
586  rec.tpc_status = state.GetStatusCode(); /* relay status received from remote side to the client */
587  ss << ": remote host returned '" << rec.tpc_status << " "
588  << httpStatusToString(rec.tpc_status) << "' while fetching file size";
589  break;
590  case CURLE_COULDNT_CONNECT: /* socket callback may have failed */
591  switch (rec.tpc_status) {
592  case 403:
593  ss << ": connection to local/private addresses is forbidden";
594  break;
595  default:
596  ss << ": internal server failure";
597  rec.tpc_status = 500;
598  }
599  break;
600  default:
601  rec.tpc_status = 500;
602  state.SetErrorCode(500);
603  }
604  }
605 
606  if (rec.tpc_status >= 400) {
607  logTransferEvent(LogMask::Error, rec, "SIZE_FAIL", ss.str());
608  return shouldReturnErrorToClient ? req.SendSimpleResp(rec.tpc_status, NULL, NULL, generateClientErr(ss, rec, res).c_str(), 0) : -1;
609  }
610 
611  success = true;
612  ss << "Successfully determined remote size for pull request: " << state.GetContentLength();
613  logTransferEvent(LogMask::Debug, rec, "SIZE_SUCCESS", ss.str());
614  return 0;
615 }
616 
617 int TPCHandler::GetContentLengthTPCPull(CURL *curl, XrdHttpExtReq &req, uint64_t &contentLength, bool & success, TPCLogRecord &rec) {
618  State state(curl,req.tpcForwardCreds);
619  //Don't forget to copy the headers of the client's request before doing the HEAD call. Otherwise, if there is a need for authentication,
620  //it will fail
621  state.SetupHeaders(req);
622  int result;
623  //In case we cannot get the content length, we return the error to the client
624  if ((result = DetermineXferSize(curl, req, state, success, rec)) || !success) {
625  return result;
626  }
627  contentLength = state.GetContentLength();
628  return result;
629 }
630 
631 /******************************************************************************/
632 /* T P C H a n d l e r : : S e n d P e r f M a r k e r */
633 /******************************************************************************/
634 
635 int TPCHandler::SendPerfMarker(XrdHttpExtReq &req, TPCLogRecord &rec, TPC::State &state) {
636  std::stringstream ss;
637  const std::string crlf = "\n";
638  ss << "Perf Marker" << crlf;
639  ss << "Timestamp: " << time(NULL) << crlf;
640  ss << "Stripe Index: 0" << crlf;
641  ss << "Stripe Bytes Transferred: " << state.BytesTransferred() << crlf;
642  ss << "Total Stripe Count: 1" << crlf;
643  // Include the TCP connection associated with this transfer; used by
644  // the TPC client for monitoring purposes.
645  std::string desc = state.GetConnectionDescription();
646  if (!desc.empty())
647  ss << "RemoteConnections: " << desc << crlf;
648  ss << "End" << crlf;
649  rec.bytes_transferred = state.BytesTransferred();
650  logTransferEvent(LogMask::Debug, rec, "PERF_MARKER");
651 
652  return req.ChunkResp(ss.str().c_str(), 0);
653 }
654 
655 /******************************************************************************/
656 /* T P C H a n d l e r : : S e n d P e r f M a r k e r */
657 /******************************************************************************/
658 
659 int TPCHandler::SendPerfMarker(XrdHttpExtReq &req, TPCLogRecord &rec, std::vector<State*> &state,
660  off_t bytes_transferred)
661 {
662  // The 'performance marker' format is largely derived from how GridFTP works
663  // (e.g., the concept of `Stripe` is not quite so relevant here). See:
664  // https://twiki.cern.ch/twiki/bin/view/LCG/HttpTpcTechnical
665  // Example marker:
666  // Perf Marker\n
667  // Timestamp: 1537788010\n
668  // Stripe Index: 0\n
669  // Stripe Bytes Transferred: 238745\n
670  // Total Stripe Count: 1\n
671  // RemoteConnections: tcp:129.93.3.4:1234,tcp:[2600:900:6:1301:268a:7ff:fef6:a590]:2345\n
672  // End\n
673  //
674  std::stringstream ss;
675  const std::string crlf = "\n";
676  ss << "Perf Marker" << crlf;
677  ss << "Timestamp: " << time(NULL) << crlf;
678  ss << "Stripe Index: 0" << crlf;
679  ss << "Stripe Bytes Transferred: " << bytes_transferred << crlf;
680  ss << "Total Stripe Count: 1" << crlf;
681  // Build a list of TCP connections associated with this transfer; used by
682  // the TPC client for monitoring purposes.
683  bool first = true;
684  std::stringstream ss2;
685  for (std::vector<State*>::const_iterator iter = state.begin();
686  iter != state.end(); iter++)
687  {
688  std::string desc = (*iter)->GetConnectionDescription();
689  if (!desc.empty()) {
690  ss2 << (first ? "" : ",") << desc;
691  first = false;
692  }
693  }
694  if (!first)
695  ss << "RemoteConnections: " << ss2.str() << crlf;
696  ss << "End" << crlf;
697  rec.bytes_transferred = bytes_transferred;
698  logTransferEvent(LogMask::Debug, rec, "PERF_MARKER");
699 
700  return req.ChunkResp(ss.str().c_str(), 0);
701 }
702 
703 /******************************************************************************/
704 /* T P C H a n d l e r : : R u n C u r l W i t h U p d a t e s */
705 /******************************************************************************/
706 
707 int TPCHandler::RunCurlWithUpdates(CURL *curl, XrdHttpExtReq &req, State &state,
708  TPCLogRecord &rec)
709 {
710  // Create the multi-handle and add in the current transfer to it.
711  CURLM *multi_handle = curl_multi_init();
712  if (!multi_handle) {
713  rec.status = 500;
714  logTransferEvent(LogMask::Error, rec, "CURL_INIT_FAIL",
715  "Failed to initialize a libcurl multi-handle");
716  std::stringstream ss;
717  ss << "Failed to initialize internal server memory";
718  return req.SendSimpleResp(rec.status, NULL, NULL, generateClientErr(ss, rec).c_str(), 0);
719  }
720 
721  //curl_easy_setopt(curl, CURLOPT_BUFFERSIZE, 128*1024);
722 
723  CURLMcode mres;
724  mres = curl_multi_add_handle(multi_handle, curl);
725  if (mres) {
726  rec.status = 500;
727  std::stringstream ss;
728  ss << "Failed to add transfer to libcurl multi-handle: HTTP library failure=" << curl_multi_strerror(mres);
729  logTransferEvent(LogMask::Error, rec, "CURL_INIT_FAIL", ss.str());
730  curl_multi_cleanup(multi_handle);
731  return req.SendSimpleResp(rec.status, NULL, NULL, generateClientErr(ss, rec).c_str(), 0);
732  }
733 
734  // Start response to client prior to the first call to curl_multi_perform
735  int retval = req.StartChunkedResp(201, "Created", "Content-Type: text/plain");
736  if (retval) {
737  curl_multi_cleanup(multi_handle);
738  logTransferEvent(LogMask::Error, rec, "RESPONSE_FAIL",
739  "Failed to send the initial response to the TPC client");
740  return retval;
741  } else {
742  logTransferEvent(LogMask::Debug, rec, "RESPONSE_START",
743  "Initial transfer response sent to the TPC client");
744  }
745 
746  // Transfer loop: use curl to actually run the transfer, but periodically
747  // interrupt things to send back performance updates to the client.
748  int running_handles = 1;
749  time_t last_marker = 0;
750  // Track how long it's been since the last time we recorded more bytes being transferred.
751  off_t last_advance_bytes = 0;
752  time_t last_advance_time = time(NULL);
753  time_t transfer_start = last_advance_time;
754  CURLcode res = static_cast<CURLcode>(-1);
755  do {
756  time_t now = time(NULL);
757  time_t next_marker = last_marker + m_marker_period;
758  if (now >= next_marker) {
759  off_t bytes_xfer = state.BytesTransferred();
760  if (bytes_xfer > last_advance_bytes) {
761  last_advance_bytes = bytes_xfer;
762  last_advance_time = now;
763  }
764  if (SendPerfMarker(req, rec, state)) {
765  curl_multi_remove_handle(multi_handle, curl);
766  curl_multi_cleanup(multi_handle);
767  logTransferEvent(LogMask::Error, rec, "PERFMARKER_FAIL",
768  "Failed to send a perf marker to the TPC client");
769  return -1;
770  }
771  int timeout = (transfer_start == last_advance_time) ? m_first_timeout : m_timeout;
772  if (now > last_advance_time + timeout) {
773  const char *log_prefix = rec.log_prefix.c_str();
774  bool tpc_pull = strncmp("Pull", log_prefix, 4) == 0;
775 
777  std::stringstream ss;
778  ss << "Transfer failed because no bytes have been "
779  << (tpc_pull ? "received from the source (pull mode) in "
780  : "transmitted to the destination (push mode) in ") << timeout << " seconds.";
781  state.SetErrorMessage(ss.str());
782  curl_multi_remove_handle(multi_handle, curl);
783  curl_multi_cleanup(multi_handle);
784  break;
785  }
786  last_marker = now;
787  }
788  // The transfer will start after this point, notify the packet marking manager
789  rec.pmarkManager.startTransfer();
790  mres = curl_multi_perform(multi_handle, &running_handles);
791  if (mres == CURLM_CALL_MULTI_PERFORM) {
792  // curl_multi_perform should be called again immediately. On newer
793  // versions of curl, this is no longer used.
794  continue;
795  } else if (mres != CURLM_OK) {
796  break;
797  } else if (running_handles == 0) {
798  break;
799  }
800 
801  rec.pmarkManager.beginPMarks();
802  //printf("There are %d running handles\n", running_handles);
803 
804  // Harvest any messages, looking for CURLMSG_DONE.
805  CURLMsg *msg;
806  do {
807  int msgq = 0;
808  msg = curl_multi_info_read(multi_handle, &msgq);
809  if (msg && (msg->msg == CURLMSG_DONE)) {
810  CURL *easy_handle = msg->easy_handle;
811  res = msg->data.result;
812  curl_multi_remove_handle(multi_handle, easy_handle);
813  }
814  } while (msg);
815 
816  int64_t max_sleep_time = next_marker - time(NULL);
817  if (max_sleep_time <= 0) {
818  continue;
819  }
820  int fd_count;
821  mres = curl_multi_wait(multi_handle, NULL, 0, max_sleep_time*1000, &fd_count);
822  if (mres != CURLM_OK) {
823  break;
824  }
825  } while (running_handles);
826 
827  if (mres != CURLM_OK) {
828  std::stringstream ss;
829  ss << "Internal libcurl multi-handle error: HTTP library failure=" << curl_multi_strerror(mres);
830  logTransferEvent(LogMask::Error, rec, "TRANSFER_CURL_ERROR", ss.str());
831 
832  curl_multi_remove_handle(multi_handle, curl);
833  curl_multi_cleanup(multi_handle);
834 
835  if ((retval = req.ChunkResp(generateClientErr(ss, rec).c_str(), 0))) {
836  logTransferEvent(LogMask::Error, rec, "RESPONSE_FAIL",
837  "Failed to send error message to the TPC client");
838  return retval;
839  }
840  return req.ChunkResp(NULL, 0);
841  }
842 
843  // Harvest any messages, looking for CURLMSG_DONE.
844  CURLMsg *msg;
845  do {
846  int msgq = 0;
847  msg = curl_multi_info_read(multi_handle, &msgq);
848  if (msg && (msg->msg == CURLMSG_DONE)) {
849  CURL *easy_handle = msg->easy_handle;
850  res = msg->data.result;
851  curl_multi_remove_handle(multi_handle, easy_handle);
852  }
853  } while (msg);
854 
855  if (!state.GetErrorCode() && res == static_cast<CURLcode>(-1)) { // No transfers returned?!?
856  curl_multi_remove_handle(multi_handle, curl);
857  curl_multi_cleanup(multi_handle);
858  std::stringstream ss;
859  ss << "Internal state error in libcurl";
860  logTransferEvent(LogMask::Error, rec, "TRANSFER_CURL_ERROR", ss.str());
861 
862  if ((retval = req.ChunkResp(generateClientErr(ss, rec).c_str(), 0))) {
863  logTransferEvent(LogMask::Error, rec, "RESPONSE_FAIL",
864  "Failed to send error message to the TPC client");
865  return retval;
866  }
867  return req.ChunkResp(NULL, 0);
868  }
869  curl_multi_cleanup(multi_handle);
870 
871  // The transfer is over at this point: any error recorded so far - a failed
872  // write to the local filesystem or the stall detector having fired - is the
873  // reason why the transfer failed. Flushing and closing the destination file
874  // below may fail as well but, as such a failure is usually a consequence of
875  // the transfer failure, it must not be reported instead of it.
876  const int transferErrorCode = state.GetErrorCode();
877  std::string transferErrorMsg = state.GetErrorMessage();
878 
879  state.Flush();
880 
881  rec.bytes_transferred = state.BytesTransferred();
882  rec.tpc_status = state.GetStatusCode();
883 
884  // Explicitly finalize the stream (which will close the underlying file
885  // handle) before the response is sent. In some cases, subsequent HTTP
886  // requests can occur before the filesystem is done closing the handle -
887  // and those requests may occur against partial data.
888  state.Finalize();
889 
890  // A failure to flush or to close the destination file is always logged and is
891  // appended to the error reported to the client, but it never replaces the
892  // transfer failure itself: it is usually a consequence of it.
893  std::string finalizeErrorMsg, finalizeErrorSuffix;
894  if (state.GetFinalizeErrorCode()) {
895  std::stringstream ss2;
896  ss2 << (state.GetFinalizeErrorCode() == State::errFlush
897  ? "Failed to flush the file to the local filesystem."
898  : "Failed to finalize and close file handle.");
899  std::string err = state.GetFinalizeErrorMessage();
900  if (!err.empty()) {
901  std::replace(err.begin(), err.end(), '\n', ' ');
902  ss2 << " " << err;
903  }
904  finalizeErrorMsg = ss2.str();
905  logTransferEvent(LogMask::Error, rec, "CLOSE_FAIL", finalizeErrorMsg);
906  finalizeErrorSuffix = "; " + finalizeErrorMsg;
907  }
908 
909  // Generate the final response back to the client.
910  std::stringstream ss;
911  bool success = false;
912  if (state.GetStatusCode() >= 400) {
913  std::string err = state.GetErrorMessage();
914  std::stringstream ss2;
915  ss2 << "Remote side failed with status code " << state.GetStatusCode();
916  if (!err.empty()) {
917  std::replace(err.begin(), err.end(), '\n', ' ');
918  ss2 << "; error message: \"" << err << "\"";
919  }
920  logTransferEvent(LogMask::Error, rec, "TRANSFER_FAIL", ss2.str());
921  ss2 << finalizeErrorSuffix;
922  ss << generateClientErr(ss2, rec);
923  } else if (transferErrorCode == State::errTimeout) {
924  // The stall detector fired; its message already describes precisely
925  // what happened, report it as-is.
926  std::stringstream ss2;
927  ss2 << transferErrorMsg;
928  logTransferEvent(LogMask::Error, rec, "TRANSFER_FAIL", ss2.str());
929  ss2 << finalizeErrorSuffix;
930  ss << generateClientErr(ss2, rec);
931  } else if (transferErrorCode) {
932  if (transferErrorMsg.empty()) {transferErrorMsg = "(no error message provided)";}
933  else {std::replace(transferErrorMsg.begin(), transferErrorMsg.end(), '\n', ' ');}
934  std::stringstream ss2;
935  ss2 << "Error when interacting with local filesystem: " << transferErrorMsg;
936  logTransferEvent(LogMask::Error, rec, "TRANSFER_FAIL", ss2.str());
937  ss2 << finalizeErrorSuffix;
938  ss << generateClientErr(ss2, rec);
939  } else if (res != CURLE_OK) {
940  std::stringstream ss2;
941  ss2 << "Internal transfer failure";
942  std::stringstream ss3;
943  ss3 << ss2.str() << ": " << curl_easy_strerror(res);
944  logTransferEvent(LogMask::Error, rec, "TRANSFER_FAIL", ss3.str());
945  ss2 << finalizeErrorSuffix;
946  ss << generateClientErr(ss2, rec, res);
947  } else if (!finalizeErrorMsg.empty()) {
948  // Nothing else went wrong: the flush/close failure is the reason of the failure.
949  std::stringstream ss2;
950  ss2 << finalizeErrorMsg;
951  ss << generateClientErr(ss2, rec);
952  } else {
953  ss << "success: Created";
954  success = true;
955  }
956 
957  if ((retval = req.ChunkResp(ss.str().c_str(), 0))) {
958  logTransferEvent(LogMask::Error, rec, "TRANSFER_ERROR",
959  "Failed to send last update to remote client");
960  return retval;
961  } else if (success) {
962  logTransferEvent(LogMask::Info, rec, "TRANSFER_SUCCESS");
963  rec.status = 0;
964  }
965  return req.ChunkResp(NULL, 0);
966 }
967 
968 /******************************************************************************/
969 /* T P C H a n d l e r : : P r o c e s s P u s h R e q */
970 /******************************************************************************/
971 
972 int TPCHandler::ProcessPushReq(const std::string & resource, XrdHttpExtReq &req) {
973  TPCLogRecord rec(req, TpcType::Push);
974  rec.allow_local = m_allow_local;
975  rec.allow_private = m_allow_private;
976  rec.log_prefix = "PushRequest";
977  rec.local = req.resource;
978  rec.remote = resource;
979  rec.m_log = &m_log;
980  char *name = req.GetSecEntity().name;
981  req.GetClientID(rec.clID);
982  if (name) rec.name = name;
983  logTransferEvent(LogMask::Info, rec, "PUSH_START", "Starting a push request");
984 
985  ManagedCurlHandle curlPtr(curl_easy_init());
986  auto curl = curlPtr.get();
987  if (!curl) {
988  std::stringstream ss;
989  ss << "Failed to initialize internal transfer resources";
990  rec.status = 500;
991  logTransferEvent(LogMask::Error, rec, "PUSH_FAIL", ss.str());
992  return req.SendSimpleResp(rec.status, NULL, NULL, generateClientErr(ss, rec).c_str(), 0);
993  }
994  ConfigureCurlLowSpeed(curl);
995  curl_easy_setopt(curl, CURLOPT_NOSIGNAL, 1);
996  curl_easy_setopt(curl, CURLOPT_SSLVERSION, CURL_SSLVERSION_TLSv1_2);
997  curl_easy_setopt(curl, CURLOPT_HTTP_VERSION, (long) CURL_HTTP_VERSION_1_1);
998 #if CURL_AT_LEAST_VERSION(7, 85, 0)
999  curl_easy_setopt(curl, CURLOPT_PROTOCOLS_STR, "https,http");
1000  curl_easy_setopt(curl, CURLOPT_REDIR_PROTOCOLS_STR, "https,http");
1001 #else
1002  long protocols = CURLPROTO_HTTP | CURLPROTO_HTTPS;
1003  curl_easy_setopt(curl, CURLOPT_PROTOCOLS, protocols);
1004  curl_easy_setopt(curl, CURLOPT_REDIR_PROTOCOLS, protocols);
1005 #endif
1006  curl_easy_setopt(curl, CURLOPT_OPENSOCKETFUNCTION, opensocket_callback);
1007  curl_easy_setopt(curl, CURLOPT_OPENSOCKETDATA, &rec);
1008  curl_easy_setopt(curl, CURLOPT_CLOSESOCKETFUNCTION, closesocket_callback);
1009  curl_easy_setopt(curl, CURLOPT_SOCKOPTFUNCTION, sockopt_callback);
1010  curl_easy_setopt(curl, CURLOPT_CLOSESOCKETDATA, &rec);
1011  curl_easy_setopt(curl, CURLOPT_CONNECTTIMEOUT, CONNECT_TIMEOUT);
1012 
1013  auto query_header = XrdOucTUtils::caseInsensitiveFind(req.headers,"xrd-http-fullresource");
1014  std::string redirect_resource = req.resource;
1015  if (query_header != req.headers.end()) {
1016  redirect_resource = query_header->second;
1017  }
1018 
1019  AtomicBeg(m_monid_mutex);
1020  uint64_t file_monid = AtomicInc(m_monid);
1021  AtomicEnd(m_monid_mutex);
1022  std::unique_ptr<XrdSfsFile> fh(m_sfs->newFile(name, file_monid));
1023  if (!fh.get()) {
1024  rec.status = 500;
1025  std::stringstream ss;
1026  ss << "Failed to initialize internal transfer file handle";
1027  logTransferEvent(LogMask::Error, rec, "OPEN_FAIL",
1028  ss.str());
1029  return req.SendSimpleResp(rec.status, NULL, NULL, generateClientErr(ss, rec).c_str(), 0);
1030  }
1031  std::string full_url = prepareURL(req);
1032 
1033  std::string authz = GetAuthz(req);
1034 
1035  int open_results = OpenWaitStall(*fh, full_url, SFS_O_RDONLY, 0644,
1036  req.GetSecEntity(), authz);
1037  if (SFS_REDIRECT == open_results) {
1038  int result = RedirectTransfer(curl, redirect_resource, req, fh->error, rec);
1039  return result;
1040  } else if (SFS_OK != open_results) {
1041  int code;
1042  std::stringstream ss;
1043  const char *msg = fh->error.getErrText(code);
1044  if (msg == NULL) ss << "Failed to open local resource";
1045  else ss << msg;
1046  rec.status = mapErrNoToHttp(code);
1047  logTransferEvent(LogMask::Error, rec, "OPEN_FAIL", msg);
1048  int resp_result = req.SendSimpleResp(rec.status, NULL, NULL, generateClientErr(ss, rec).c_str(), 0);
1049  fh->close();
1050  return resp_result;
1051  }
1052  if (!ConfigureCurlCA(curl, rec)) {
1053  std::stringstream ss;
1054  ss << "Failed to configure the certificate authorities for the transfer";
1055  rec.status = 500;
1056  logTransferEvent(LogMask::Error, rec, "PUSH_FAIL", ss.str());
1057  int resp_result = req.SendSimpleResp(rec.status, NULL, NULL, generateClientErr(ss, rec).c_str(), 0);
1058  fh->close();
1059  return resp_result;
1060  }
1061  curl_easy_setopt(curl, CURLOPT_URL, resource.c_str());
1062 
1063  Stream stream(std::move(fh), 0, 0, m_log);
1064  State state(0, stream, curl, true, req.tpcForwardCreds);
1065  state.SetupHeaders(req);
1066 
1067  return RunCurlWithUpdates(curl, req, state, rec);
1068 }
1069 
1070 /******************************************************************************/
1071 /* T P C H a n d l e r : : P r o c e s s P u l l R e q */
1072 /******************************************************************************/
1073 
1074 int TPCHandler::ProcessPullReq(const std::string &resource, XrdHttpExtReq &req) {
1075  TPCLogRecord rec(req,TpcType::Pull);
1076  rec.allow_local = m_allow_local;
1077  rec.allow_private = m_allow_private;
1078  rec.log_prefix = "PullRequest";
1079  rec.local = req.resource;
1080  rec.remote = resource;
1081  rec.m_log = &m_log;
1082  char *name = req.GetSecEntity().name;
1083  req.GetClientID(rec.clID);
1084  if (name) rec.name = name;
1085  logTransferEvent(LogMask::Info, rec, "PULL_START", "Starting a pull request");
1086 
1087  ManagedCurlHandle curlPtr(curl_easy_init());
1088  auto curl = curlPtr.get();
1089  if (!curl) {
1090  std::stringstream ss;
1091  ss << "Failed to initialize internal transfer resources";
1092  rec.status = 500;
1093  logTransferEvent(LogMask::Error, rec, "PULL_FAIL", ss.str());
1094  return req.SendSimpleResp(rec.status, NULL, NULL, generateClientErr(ss, rec).c_str(), 0);
1095  }
1096  ConfigureCurlLowSpeed(curl);
1097 
1098  // ddavila 2023-01-05:
1099  // The following change was required by the Rucio/SENSE project where
1100  // multiple IP addresses, each from a different subnet, are assigned to a
1101  // single server and routed differently by SENSE.
1102  // The above requires the server to utilize the same IP, that was used to
1103  // start the TPC, for the resolution of the given TPC instead of
1104  // using any of the IPs available.
1105  if (m_fixed_route) {
1106  // Get the hostname used to contact the server from the http header
1107  std::string host;
1108  auto host_header = XrdOucTUtils::caseInsensitiveFind(req.headers, "host");
1109 
1110  if (host_header != req.headers.end()) {
1111  host = host_header->second;
1112  }
1113 
1114  // Get the IP addresses associated with the hostname
1115  char ip[64]; // IPv6 addresses are up to 45 characters long
1116  std::vector<XrdNetAddr> addresses;
1117  const char *eText = XrdNetUtils::GetAddrs(host, addresses, nullptr, XrdNetUtils::prefAuto, 0);
1118 
1119  if (eText || addresses.empty() ||
1120  addresses.front().Format(ip, sizeof(ip), XrdNetAddrInfo::fmtAddr,XrdNetAddrInfo::noPortRaw) <= 0) {
1121  std::stringstream ss;
1122  ss << "Failed to determine host address of incoming request";
1123  rec.status = 500;
1124  logTransferEvent(LogMask::Error, rec, "PULL_FAIL", ss.str());
1125  return req.SendSimpleResp(rec.status, NULL, NULL, generateClientErr(ss, rec).c_str(), 0);
1126  }
1127 
1128  logTransferEvent(LogMask::Info, rec, "LOCAL IP", ip);
1129  curl_easy_setopt(curl, CURLOPT_INTERFACE, ip);
1130  }
1131  curl_easy_setopt(curl, CURLOPT_NOSIGNAL, 1);
1132  curl_easy_setopt(curl, CURLOPT_SSLVERSION, CURL_SSLVERSION_TLSv1_2);
1133  curl_easy_setopt(curl, CURLOPT_HTTP_VERSION, (long) CURL_HTTP_VERSION_1_1);
1134 #if CURL_AT_LEAST_VERSION(7, 85, 0)
1135  curl_easy_setopt(curl, CURLOPT_PROTOCOLS_STR, "https,http");
1136  curl_easy_setopt(curl, CURLOPT_REDIR_PROTOCOLS_STR, "https,http");
1137 #else
1138  long protocols = CURLPROTO_HTTP | CURLPROTO_HTTPS;
1139  curl_easy_setopt(curl, CURLOPT_PROTOCOLS, protocols);
1140  curl_easy_setopt(curl, CURLOPT_REDIR_PROTOCOLS, protocols);
1141 #endif
1142  curl_easy_setopt(curl, CURLOPT_OPENSOCKETFUNCTION, opensocket_callback);
1143  curl_easy_setopt(curl, CURLOPT_OPENSOCKETDATA, &rec);
1144  curl_easy_setopt(curl, CURLOPT_SOCKOPTFUNCTION, sockopt_callback);
1145  curl_easy_setopt(curl, CURLOPT_SOCKOPTDATA , &rec);
1146  curl_easy_setopt(curl, CURLOPT_CLOSESOCKETFUNCTION, closesocket_callback);
1147  curl_easy_setopt(curl, CURLOPT_CLOSESOCKETDATA, &rec);
1148  curl_easy_setopt(curl, CURLOPT_CONNECTTIMEOUT, CONNECT_TIMEOUT);
1149  std::unique_ptr<XrdSfsFile> fh(m_sfs->newFile(name, m_monid++));
1150  if (!fh.get()) {
1151  std::stringstream ss;
1152  ss << "Failed to initialize internal transfer file handle";
1153  rec.status = 500;
1154  logTransferEvent(LogMask::Error, rec, "PULL_FAIL", ss.str());
1155  return req.SendSimpleResp(rec.status, NULL, NULL, generateClientErr(ss, rec).c_str(), 0);
1156  }
1157  auto query_header = XrdOucTUtils::caseInsensitiveFind(req.headers,"xrd-http-fullresource");
1158  std::string redirect_resource = req.resource;
1159  if (query_header != req.headers.end()) {
1160  redirect_resource = query_header->second;
1161  }
1163  auto overwrite_header = XrdOucTUtils::caseInsensitiveFind(req.headers,"overwrite");
1164  if ((overwrite_header == req.headers.end()) || (overwrite_header->second == "T")) {
1165  if (! usingEC) mode = SFS_O_TRUNC;
1166  }
1167  int streams = 1;
1168  {
1169  auto streams_header = XrdOucTUtils::caseInsensitiveFind(req.headers,"x-number-of-streams");
1170  if (streams_header != req.headers.end()) {
1171  int stream_req = -1;
1172  try {
1173  stream_req = std::stol(streams_header->second);
1174  } catch (...) { // Handled below
1175  }
1176  if (stream_req < 0 || stream_req > 100) {
1177  std::stringstream ss;
1178  ss << "Invalid request for number of streams";
1179  rec.status = 400;
1180  logTransferEvent(LogMask::Info, rec, "INVALID_REQUEST", ss.str());
1181  return req.SendSimpleResp(rec.status, NULL, NULL, generateClientErr(ss, rec).c_str(), 0);
1182  }
1183  streams = stream_req == 0 ? 1 : stream_req;
1184  }
1185  }
1186  rec.streams = streams;
1187  std::string full_url = prepareURL(req);
1188  std::string authz = GetAuthz(req);
1189  curl_easy_setopt(curl, CURLOPT_URL, resource.c_str());
1190  if (!ConfigureCurlCA(curl, rec)) {
1191  std::stringstream ss;
1192  ss << "Failed to configure the certificate authorities for the transfer";
1193  rec.status = 500;
1194  logTransferEvent(LogMask::Error, rec, "PULL_FAIL", ss.str());
1195  return req.SendSimpleResp(rec.status, NULL, NULL, generateClientErr(ss, rec).c_str(), 0);
1196  }
1197  uint64_t sourceFileContentLength = 0;
1198  {
1199  //Get the content-length of the source file and pass it to the OSS layer
1200  //during the open
1201  bool success;
1202  GetContentLengthTPCPull(curl, req, sourceFileContentLength, success, rec);
1203  if(success) {
1204  //In the case we cannot get the information from the source server (offline or other error)
1205  //we just don't add the size information to the opaque of the local file to open
1206  full_url += "&oss.asize=" + std::to_string(sourceFileContentLength);
1207  } else {
1208  // In the case the GetContentLength is not successful, an error will be returned to the client
1209  // just exit here so we don't open the file!
1210  return 0;
1211  }
1212  }
1213  int open_result = OpenWaitStall(*fh, full_url, mode|SFS_O_WRONLY,
1214  0644 | SFS_O_MKPTH,
1215  req.GetSecEntity(), authz);
1216  if (SFS_REDIRECT == open_result) {
1217  int result = RedirectTransfer(curl, redirect_resource, req, fh->error, rec);
1218  return result;
1219  } else if (SFS_OK != open_result) {
1220  int code;
1221  std::stringstream ss;
1222  const char *msg = fh->error.getErrText(code);
1223  if ((msg == NULL) || (*msg == '\0')) ss << "Failed to open local resource";
1224  else ss << msg;
1225  rec.status = mapErrNoToHttp(code);
1226  logTransferEvent(LogMask::Error, rec, "OPEN_FAIL", ss.str());
1227  int resp_result = req.SendSimpleResp(rec.status, NULL, NULL,
1228  generateClientErr(ss, rec).c_str(), 0);
1229  fh->close();
1230  return resp_result;
1231  }
1232  Stream stream(std::move(fh), streams * m_pipelining_multiplier, streams > 1 ? m_block_size : m_small_block_size, m_log);
1233  State state(0, stream, curl, false, req.tpcForwardCreds);
1234  state.SetupHeaders(req);
1235  state.SetContentLength(sourceFileContentLength);
1236 
1237  if (streams > 1) {
1238  return RunCurlWithStreams(req, state, streams, rec);
1239  } else {
1240  return RunCurlWithUpdates(curl, req, state, rec);
1241  }
1242 }
1243 
1244 /******************************************************************************/
1245 /* T P C H a n d l e r : : l o g T r a n s f e r E v e n t */
1246 /******************************************************************************/
1247 
1248 void TPCHandler::logTransferEvent(LogMask mask, const TPCLogRecord &rec,
1249  const std::string &event, const std::string &message)
1250 {
1251  if (!(m_log.getMsgMask() & mask)) {return;}
1252 
1253  std::stringstream ss;
1254  ss << "event=" << event << ", local=" << rec.local << ", remote=" << rec.remote;
1255  if (rec.name.empty())
1256  ss << ", user=(anonymous)";
1257  else
1258  ss << ", user=" << rec.name;
1259  if (rec.streams != 1)
1260  ss << ", streams=" << rec.streams;
1261  if (rec.bytes_transferred >= 0)
1262  ss << ", bytes_transferred=" << rec.bytes_transferred;
1263  if (rec.status >= 0)
1264  ss << ", status=" << rec.status;
1265  if (rec.tpc_status >= 0)
1266  ss << ", tpc_status=" << rec.tpc_status;
1267  if (!message.empty())
1268  ss << "; " << message;
1269  m_log.Log(mask, rec.log_prefix.c_str(), ss.str().c_str());
1270 }
1271 
1272 std::string TPCHandler::generateClientErr(std::stringstream &err_ss, const TPCLogRecord &rec, CURLcode cCode) {
1273  std::stringstream ssret;
1274  ssret << "failure: " << err_ss.str() << ", local=" << rec.local <<", remote=" << rec.remote;
1275  if(cCode != CURLcode::CURLE_OK) {
1276  ssret << ", HTTP library failure=" << curl_easy_strerror(cCode);
1277  }
1278  return ssret.str();
1279 }
1280 /******************************************************************************/
1281 /* X r d H t t p G e t E x t H a n d l e r */
1282 /******************************************************************************/
1283 
1284 extern "C" {
1285 
1286 XrdHttpExtHandler *XrdHttpGetExtHandler(XrdSysError *log, const char * config, const char * /*parms*/, XrdOucEnv *myEnv) {
1287  if (curl_global_init(CURL_GLOBAL_DEFAULT)) {
1288  log->Emsg("TPCInitialize", "libcurl failed to initialize");
1289  return NULL;
1290  }
1291 
1292  TPCHandler *retval{NULL};
1293  if (!config) {
1294  log->Emsg("TPCInitialize", "TPC handler requires a config filename in order to load");
1295  return NULL;
1296  }
1297  try {
1298  log->Emsg("TPCInitialize", "Will load configuration for the TPC handler from", config);
1299  retval = new TPCHandler(log, config, myEnv);
1300  } catch (std::runtime_error &re) {
1301  log->Emsg("TPCInitialize", "Encountered a runtime failure when loading ", re.what());
1302  //printf("Provided env vars: %p, XrdInet*: %p\n", myEnv, myEnv->GetPtr("XrdInet*"));
1303  }
1304  return retval;
1305 }
1306 
1307 }
void CURL
XrdVERSIONINFO(XrdHttpGetExtHandler, HttpTPC)
XrdHttpExtHandler * XrdHttpGetExtHandler(XrdSysError *log, const char *config, const char *, XrdOucEnv *myEnv)
static std::string PrepareURL(const std::string &url)
std::string encode_xrootd_opaque_to_uri(CURL *curl, const std::string &opaque)
static bool IsAllowedScheme(const std::string &url)
int mapErrNoToHttp(int errNo)
std::string httpStatusToString(int status)
Utility functions for XrdHTTP.
std::string encode_str(const std::string &str)
#define close(a)
Definition: XrdPosix.hh:48
bool Debug
void getline(uchar *buff, int blen)
#define SFS_REDIRECT
#define SFS_O_MKPTH
#define SFS_STALL
#define SFS_O_RDONLY
#define SFS_STARTED
#define SFS_O_WRONLY
#define SFS_O_CREAT
int XrdSfsFileOpenMode
#define SFS_OK
#define SFS_O_TRUNC
#define AtomicInc(x)
#define AtomicBeg(Mtx)
#define AtomicEnd(Mtx)
@ Error
int GetFinalizeErrorCode() const
int GetStatusCode() const
off_t BytesTransferred() const
void SetErrorMessage(const std::string &error_msg)
int GetErrorCode() const
std::string GetFinalizeErrorMessage() const
std::string GetErrorMessage() const
std::string GetConnectionDescription()
void SetupHeaders(XrdHttpExtReq &req)
void SetContentLength(const off_t content_length)
off_t GetContentLength() const
void SetErrorCode(int error_code)
TPCHandler(XrdSysError *log, const char *config, XrdOucEnv *myEnv)
virtual int ProcessReq(XrdHttpExtReq &req)
virtual ~TPCHandler()
virtual bool MatchesPath(const char *verb, const char *path)
Tells if the incoming path is recognized as one of the paths that have to be processed.
int ChunkResp(const char *body, long long bodylen)
Send a (potentially partial) body in a chunked response; invoking with NULL body.
void GetClientID(std::string &clid)
std::map< std::string, std::string > & headers
std::string resource
std::string verb
int StartChunkedResp(int code, const char *desc, const char *header_to_add)
Starts a chunked response; body of request is sent over multiple parts using the SendChunkResp.
const XrdSecEntity & GetSecEntity() const
int SendSimpleResp(int code, const char *desc, const char *header_to_add, const char *body, long long bodylen)
Sends a basic response. If the length is < 0 then it is calculated internally.
static std::string prepareOpenURL(const std::string &reqResource, std::map< std::string, std::string > &reqHeaders, const std::map< std::string, std::string > &hdr2cgimap)
static const int noPortRaw
Use raw address format (no port)
@ fmtAddr
Address using suitable ipv4 or ipv6 format.
static const char * GetAddrs(const char *hSpec, XrdNetAddr *aListP[], int &aListN, AddrOpts opts=allIPMap, int pNum=PortInSpec)
Definition: XrdNetUtils.cc:274
void * GetPtr(const char *varname)
Definition: XrdOucEnv.cc:281
const char * getErrText()
void setUCap(int ucval)
Set user capabilties.
static std::map< std::string, T >::const_iterator caseInsensitiveFind(const std::map< std::string, T > &m, const std::string &lowerCaseSearchKey)
Definition: XrdOucTUtils.hh:79
char * name
Entity's name.
Definition: XrdSecEntity.hh:69
virtual XrdSfsFile * newFile(char *user=0, int MonID=0)=0
XrdOucErrInfo & error
virtual int open(const char *fileName, XrdSfsFileOpenMode openMode, mode_t createMode, const XrdSecEntity *client=0, const char *opaque=0)=0
virtual int close()=0
int Emsg(const char *esfx, int ecode, const char *text1, const char *text2=0)
Definition: XrdSysError.cc:95
XrdSysLogger * logger(XrdSysLogger *lp=0)
Definition: XrdSysError.hh:141
int getMsgMask()
Definition: XrdSysError.hh:156
void Log(int mask, const char *esfx, const char *text1, const char *text2=0, const char *text3=0)
Definition: XrdSysError.hh:133
std::unique_ptr< CURL, CurlDeleter > ManagedCurlHandle
@ Warning
void operator()(CURL *curl)
static const int uIPv64
ucap: Supports only IPv4 info
static const int isaPush