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