9#include "XrdVersion.hh"
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;
43bool TPCHandler::allowMissingCRL =
false;
51TPCHandler::TPCLogRecord::~TPCLogRecord()
58 monInfo.
clID = clID.c_str();
60 gettimeofday(&monInfo.
endT, 0);
63 {monInfo.
dstURL = local.c_str();
64 monInfo.
srcURL = remote.c_str();
66 monInfo.
dstURL = remote.c_str();
67 monInfo.
srcURL = local.c_str();
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);
78 tpcMonitor->Report(monInfo);
88 if (curl) curl_easy_cleanup(curl);
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()) {
108 return CURL_SOCKOPT_ALREADY_CONNECTED;
110 return CURL_SOCKOPT_OK;
122int TPCHandler::opensocket_callback(
void *clientp,
123 curlsocktype purpose,
124 struct curl_sockaddr *aInfo)
128 if (purpose != CURLSOCKTYPE_IPCXN)
129 return CURL_SOCKET_BAD;
132 return CURL_SOCKET_BAD;
135 int fd = XrdSysFD_Socket(aInfo->family, aInfo->socktype, aInfo->protocol);
138 return CURL_SOCKET_BAD;
144 XrdNetAddr thePeer(&(aInfo->addr));
145 TPCLogRecord *rec =
static_cast<TPCLogRecord*
>(clientp);
148 if ((!rec->allow_private && thePeer.isPrivate()) || (!rec->allow_local && thePeer.isLocal())) {
149 rec->tpc_status = 403;
150 rec->m_log->Emsg(rec->log_prefix.c_str(),
151 "Connection to local/private address is forbidden");
153 return CURL_SOCKET_BAD;
158 std::stringstream connectErrMsg;
159 if(!rec->pmarkManager.connect(fd, &(aInfo->addr), aInfo->addrlen, CONNECT_TIMEOUT, connectErrMsg)) {
161 rec->m_log->Emsg(rec->log_prefix.c_str(),
"Unable to connect socket: ", connectErrMsg.str().c_str());
162 return CURL_SOCKET_BAD;
168int TPCHandler::closesocket_callback(
void *clientp, curl_socket_t fd) {
169 TPCLogRecord * rec = (TPCLogRecord *)clientp;
174 rec->pmarkManager.endPmark(fd);
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);
193 if (rec && rec->ca_store) {
197 SSL_CTX_set1_cert_store(ctx, rec->ca_store.get());
199 if (allowMissingCRL) {
203 SSL_CTX_set_verify(ctx, SSL_VERIFY_PEER, verify_callback);
208int TPCHandler::verify_callback(
int preverify_ok, X509_STORE_CTX* ctx) {
209 if (preverify_ok == 1)
return 1;
211 int err = X509_STORE_CTX_get_error(ctx);
213 if (err == X509_V_ERR_UNABLE_TO_GET_CRL) {
214 X509_STORE_CTX_set_error(ctx, X509_V_OK);
226std::string TPCHandler::prepareURL(XrdHttpExtReq &req) {
231bool TPCHandler::mismatchReprDigest(
const std::map<std::string, std::string> & passiveSrvReprDigest, XrdHttpExtReq &req,
233 if(passiveSrvReprDigest.size()) {
234 for (
const auto & [digestName, digestValue]: passiveSrvReprDigest) {
235 auto clientDigestMatch = req.
mReprDigest.find(digestName);
238 if (clientDigestMatch->second != digestValue) {
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());
246 req.
SendSimpleResp(rec.status, NULL, NULL, generateClientErr(errMsg, rec, CURLcode::CURLE_OK).c_str(), 0);
266 std::stringstream parser(opaque);
267 std::string sequence;
268 std::stringstream output;
270 while (
getline(parser, sequence,
'&')) {
271 if (sequence.empty()) {
continue;}
272 size_t equal_pos = sequence.find(
'=');
274 if (equal_pos != std::string::npos)
275 val = curl_easy_escape(curl, sequence.c_str() + equal_pos + 1, sequence.size() - equal_pos - 1);
277 if (!val && equal_pos != std::string::npos) {
continue;}
279 if (!first) output <<
"&";
281 output << sequence.substr(0, equal_pos);
283 output <<
"=" << val;
295TPCHandler::ConfigureCurlCA(
CURL *curl, TPCLogRecord &rec)
305 if (m_ca_file && m_sslctx_supported && m_cafile.empty()) {
306 rec.ca_store = m_ca_file->CAStore();
308 m_log.Log(
Error,
"TpcHandler",
"No CA store is available; refusing to "
309 "fall back to libcurl's default CA bundle");
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);
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());
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) {
333 curl_easy_setopt(curl, CURLOPT_SSL_CTX_FUNCTION, ssl_ctx_callback);
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());
341 else if (!m_cadir.empty()) {
342 curl_easy_setopt(curl, CURLOPT_CAPATH, m_cadir.c_str());
344 if (!m_cafile.empty()) {
345 curl_easy_setopt(curl, CURLOPT_CAINFO, m_cafile.c_str());
351TPCHandler::ConfigureCurlLowSpeed(
CURL *curl)
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);
365 return !strcmp(verb,
"COPY") || !strcmp(verb,
"OPTIONS");
374 const std::string replace_schemes[] = {
"davs://",
"s3://",
"s3s://" };
376 for (
const auto& s : replace_schemes)
377 if (url.compare(0, s.size(), s) == 0)
378 return "https://" + url.substr(s.size());
385 const std::string allowed_schemes[] = {
"https://",
"http://" };
387 for (
const auto& s : allowed_schemes)
388 if (url.compare(0, s.size(), s) == 0)
399 if (req.
verb ==
"OPTIONS") {
400 return ProcessOptionsReq(req);
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);
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);
419 if (srcHeader != req.
headers.end()) {
420 std::string src =
PrepareURL(srcHeader->second);
422 const char *error_src =
"COPY rejected: disallowed scheme in source URL";
423 m_log.Emsg(
"ProcessReq", error_src, src.c_str());
426 return ProcessPullReq(src, req);
428 if (dstHeader != req.
headers.end()) {
429 const std::string& dst = dstHeader->second;
431 const char *error_dst =
"COPY rejected: disallowed scheme in destination URL";
432 m_log.Emsg(
"ProcessReq", error_dst, dst.c_str());
435 return ProcessPushReq(dst, req);
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);
454 m_allow_local(false),
455 m_allow_private(true),
457 m_fixed_route(false),
458 m_low_speed_limit(10*1024),
459 m_low_speed_time(2*60),
461 m_first_timeout(120),
462 m_log(log->logger(),
"TPC_"),
465 if (!Configure(config, myEnv)) {
466 throw std::runtime_error(
"Failed to configure the HTTP third-party-copy handler.");
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);
494 if (authz_header != req.
headers.end()) {
495 std::stringstream ss;
496 ss <<
"authz=" <<
encode_str(authz_header->second);
506int TPCHandler::RedirectTransfer(
CURL *curl,
const std::string &redirect_resource,
507 XrdHttpExtReq &req, XrdOucErrInfo &error, TPCLogRecord &rec)
511 if ((ptr == NULL) || (*ptr ==
'\0') || (port == 0)) {
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);
522 std::string finalTarget = ptr;
534 std::string newTarget;
539 finalTarget = std::move(newTarget);
541 logTransferEvent(
LogMask::Info, rec,
"REDIRECT_PLUGIN_REWRITE",
545 std::stringstream ess;
546 ess <<
"Redirect plugin error: " << errMsg;
550 generateClientErr(ess, rec).c_str(), 0);
562 std::string opaque = cgi.empty() ? std::string() : cgi.substr(1);
564 std::stringstream ss;
565 ss <<
"Location: http" << (m_desthttps ?
"s" :
"") <<
"://" << host <<
":" << port <<
"/" << redirect_resource;
567 if (!opaque.empty()) {
571 char sep = (redirect_resource.find(
'?') == std::string::npos) ?
'?' :
'&';
572 ss << sep << encode_xrootd_opaque_to_uri(curl, opaque);
577 return req.
SendSimpleResp(rec.status, NULL,
const_cast<char *
>(ss.str().c_str()),
585int TPCHandler::OpenWaitStall(XrdSfsFile &fh,
const std::string &resource,
586 int mode,
int openMode,
const XrdSecEntity &sec,
587 const std::string &authz)
594 size_t pos = resource.find(
'?');
596 std::string path = resource.substr(0, pos);
598 if (pos != std::string::npos) {
599 opaque = resource.substr(pos + 1);
604 opaque += (opaque.empty() ?
"" :
"&");
607 open_result = fh.
open(path.c_str(), mode, openMode, &sec, opaque.c_str());
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));
628int TPCHandler::PerformHEADRequest(
CURL *curl, XrdHttpExtReq &req,
State &state,
629 bool &success, TPCLogRecord &rec,
bool shouldReturnErrorToClient) {
631 curl_easy_setopt(curl, CURLOPT_NOBODY, 1);
633 curl_easy_setopt(curl, CURLOPT_TIMEOUT, CONNECT_TIMEOUT);
635 res = curl_easy_perform(curl);
638 curl_easy_setopt(curl, CURLOPT_NOBODY, 0);
640 curl_easy_setopt(curl, CURLOPT_TIMEOUT, 0L);
641 curl_easy_setopt(curl, CURLOPT_FAILONERROR,
true);
643 std::stringstream ss;
646 res = CURLE_HTTP_RETURNED_ERROR;
648 if (res != CURLE_OK) {
649 ss << curl_easy_strerror(res);
651 case CURLE_HTTP_RETURNED_ERROR:
653 ss <<
": remote host returned '" << rec.tpc_status <<
" "
656 case CURLE_COULDNT_CONNECT:
657 switch (rec.tpc_status) {
659 ss <<
": connection to local/private addresses is forbidden";
662 ss <<
": internal server failure";
663 rec.tpc_status = 500;
667 rec.tpc_status = 500;
672 if (rec.tpc_status >= 400) {
674 return shouldReturnErrorToClient ? req.
SendSimpleResp(rec.tpc_status, NULL, NULL, generateClientErr(ss, rec, res).c_str(), 0) : -1;
678 ss <<
"Successfully determined remote file information for pull request: "
681 unsigned int cksumIndex = 1;
682 for(
const auto & [cksumType,cksumValue]: state.
GetReprDigest()) {
683 ss <<
" chksum" << cksumIndex <<
"=(" << cksumType <<
"," << cksumValue <<
")";
691int TPCHandler::GetRemoteFileInfoTPCPull(
CURL *curl, XrdHttpExtReq &req, uint64_t &contentLength, std::map<std::string,std::string> & reprDigest,
bool & success, TPCLogRecord &rec) {
698 if ((result = PerformHEADRequest(curl, req, state, success, rec)) || !success) {
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;
717 ss <<
"Total Stripe Count: 1" << crlf;
722 ss <<
"RemoteConnections: " << desc << crlf;
727 return req.
ChunkResp(ss.str().c_str(), 0);
734int TPCHandler::SendPerfMarker(XrdHttpExtReq &req, TPCLogRecord &rec, std::vector<State*> &state,
735 off_t bytes_transferred)
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;
759 std::stringstream ss2;
760 for (std::vector<State*>::const_iterator iter = state.begin();
761 iter != state.end(); iter++)
763 std::string desc = (*iter)->GetConnectionDescription();
765 ss2 << (first ?
"" :
",") << desc;
770 ss <<
"RemoteConnections: " << ss2.str() << crlf;
772 rec.bytes_transferred = bytes_transferred;
775 return req.
ChunkResp(ss.str().c_str(), 0);
782int TPCHandler::RunCurlWithUpdates(
CURL *curl, XrdHttpExtReq &req,
State &state,
786 CURLM *multi_handle = curl_multi_init();
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);
799 mres = curl_multi_add_handle(multi_handle, curl);
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);
812 curl_multi_cleanup(multi_handle);
814 "Failed to send the initial response to the TPC client");
818 "Initial transfer response sent to the TPC client");
823 int running_handles = 1;
824 time_t last_marker = 0;
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);
831 time_t now = time(NULL);
832 time_t next_marker = last_marker + m_marker_period;
833 if (now >= next_marker) {
835 if (bytes_xfer > last_advance_bytes) {
836 last_advance_bytes = bytes_xfer;
837 last_advance_time = now;
839 if (SendPerfMarker(req, rec, state)) {
840 curl_multi_remove_handle(multi_handle, curl);
841 curl_multi_cleanup(multi_handle);
843 "Failed to send a perf marker to the TPC client");
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;
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.";
857 curl_multi_remove_handle(multi_handle, curl);
858 curl_multi_cleanup(multi_handle);
864 rec.pmarkManager.startTransfer();
865 mres = curl_multi_perform(multi_handle, &running_handles);
866 if (mres == CURLM_CALL_MULTI_PERFORM) {
870 }
else if (mres != CURLM_OK) {
872 }
else if (running_handles == 0) {
876 rec.pmarkManager.beginPMarks();
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);
891 int64_t max_sleep_time = next_marker - time(NULL);
892 if (max_sleep_time <= 0) {
896 mres = curl_multi_wait(multi_handle, NULL, 0, max_sleep_time*1000, &fd_count);
897 if (mres != CURLM_OK) {
900 }
while (running_handles);
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());
907 curl_multi_remove_handle(multi_handle, curl);
908 curl_multi_cleanup(multi_handle);
910 if ((retval = req.
ChunkResp(generateClientErr(ss, rec).c_str(), 0))) {
912 "Failed to send error message to the TPC client");
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);
930 if (!state.
GetErrorCode() && res ==
static_cast<CURLcode
>(-1)) {
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());
937 if ((retval = req.
ChunkResp(generateClientErr(ss, rec).c_str(), 0))) {
939 "Failed to send error message to the TPC client");
944 curl_multi_cleanup(multi_handle);
968 std::string finalizeErrorMsg, finalizeErrorSuffix;
970 std::stringstream ss2;
972 ?
"Failed to flush the file to the local filesystem."
973 :
"Failed to finalize and close file handle.");
976 std::replace(err.begin(), err.end(),
'\n',
' ');
979 finalizeErrorMsg = ss2.str();
980 logTransferEvent(
LogMask::Error, rec,
"CLOSE_FAIL", finalizeErrorMsg);
981 finalizeErrorSuffix =
"; " + finalizeErrorMsg;
985 std::stringstream ss;
986 bool success =
false;
989 std::stringstream ss2;
990 ss2 <<
"Remote side failed with status code " << state.
GetStatusCode();
992 std::replace(err.begin(), err.end(),
'\n',
' ');
993 ss2 <<
"; error message: \"" << err <<
"\"";
995 logTransferEvent(
LogMask::Error, rec,
"TRANSFER_FAIL", ss2.str());
996 ss2 << finalizeErrorSuffix;
997 ss << generateClientErr(ss2, rec);
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()) {
1024 std::stringstream ss2;
1025 ss2 << finalizeErrorMsg;
1026 ss << generateClientErr(ss2, rec);
1028 ss <<
"success: Created";
1032 if ((retval = req.
ChunkResp(ss.str().c_str(), 0))) {
1034 "Failed to send last update to remote client");
1036 }
else if (success) {
1047int TPCHandler::ProcessPushReq(
const std::string & resource, XrdHttpExtReq &req) {
1049 rec.allow_local = m_allow_local;
1050 rec.allow_private = m_allow_private;
1051 rec.log_prefix =
"PushRequest";
1053 rec.remote = resource;
1057 if (name) rec.name = name;
1058 logTransferEvent(
LogMask::Info, rec,
"PUSH_START",
"Starting a push request");
1061 auto curl = curlPtr.get();
1063 std::stringstream ss;
1064 ss <<
"Failed to initialize internal transfer resources";
1067 return req.
SendSimpleResp(rec.status, NULL, NULL, generateClientErr(ss, rec).c_str(), 0);
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");
1077 long protocols = CURLPROTO_HTTP | CURLPROTO_HTTPS;
1078 curl_easy_setopt(curl, CURLOPT_PROTOCOLS, protocols);
1079 curl_easy_setopt(curl, CURLOPT_REDIR_PROTOCOLS, protocols);
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);
1089 std::string redirect_resource = req.
resource;
1090 if (query_header != req.
headers.end()) {
1091 redirect_resource = query_header->second;
1095 uint64_t file_monid =
AtomicInc(m_monid);
1097 std::unique_ptr<XrdSfsFile> fh(m_sfs->newFile(name, file_monid));
1100 std::stringstream ss;
1101 ss <<
"Failed to initialize internal transfer file handle";
1104 return req.
SendSimpleResp(rec.status, NULL, NULL, generateClientErr(ss, rec).c_str(), 0);
1106 std::string full_url = prepareURL(req);
1108 std::string authz = GetAuthz(req);
1110 int open_results = OpenWaitStall(*fh, full_url,
SFS_O_RDONLY, 0644,
1113 int result = RedirectTransfer(curl, redirect_resource, req, fh->
error, rec);
1115 }
else if (
SFS_OK != open_results) {
1117 std::stringstream ss;
1119 if (msg == NULL) ss <<
"Failed to open local resource";
1123 int resp_result = req.
SendSimpleResp(rec.status, NULL, NULL, generateClientErr(ss, rec).c_str(), 0);
1127 if (!ConfigureCurlCA(curl, rec)) {
1128 std::stringstream ss;
1129 ss <<
"Failed to configure the certificate authorities for the transfer";
1132 int resp_result = req.
SendSimpleResp(rec.status, NULL, NULL, generateClientErr(ss, rec).c_str(), 0);
1136 curl_easy_setopt(curl, CURLOPT_URL, resource.c_str());
1138 Stream stream(std::move(fh), 0, 0, m_log);
1142 return RunCurlWithUpdates(curl, req, state, rec);
1149int TPCHandler::ProcessPullReq(
const std::string &resource, XrdHttpExtReq &req) {
1151 rec.allow_local = m_allow_local;
1152 rec.allow_private = m_allow_private;
1153 rec.log_prefix =
"PullRequest";
1155 rec.remote = resource;
1159 if (name) rec.name = name;
1160 logTransferEvent(
LogMask::Info, rec,
"PULL_START",
"Starting a pull request");
1163 auto curl = curlPtr.get();
1165 std::stringstream ss;
1166 ss <<
"Failed to initialize internal transfer resources";
1169 return req.
SendSimpleResp(rec.status, NULL, NULL, generateClientErr(ss, rec).c_str(), 0);
1171 ConfigureCurlLowSpeed(curl);
1180 if (m_fixed_route) {
1185 int sockFD = addrInfo ? addrInfo->
SockFD() : -1;
1190 logTransferEvent(
LogMask::Error, rec,
"FIXED_ROUTE_ERR",
"Failed to determine local address of incoming fixed route request");
1193 curl_easy_setopt(curl, CURLOPT_INTERFACE, ip);
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");
1203 long protocols = CURLPROTO_HTTP | CURLPROTO_HTTPS;
1204 curl_easy_setopt(curl, CURLOPT_PROTOCOLS, protocols);
1205 curl_easy_setopt(curl, CURLOPT_REDIR_PROTOCOLS, protocols);
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++));
1216 std::stringstream ss;
1217 ss <<
"Failed to initialize internal transfer file handle";
1220 return req.
SendSimpleResp(rec.status, NULL, NULL, generateClientErr(ss, rec).c_str(), 0);
1223 std::string redirect_resource = req.
resource;
1224 if (query_header != req.
headers.end()) {
1225 redirect_resource = query_header->second;
1229 if ((overwrite_header == req.
headers.end()) || (overwrite_header->second ==
"T")) {
1235 if (streams_header != req.
headers.end()) {
1236 int stream_req = -1;
1238 stream_req = std::stol(streams_header->second);
1241 if (stream_req < 0 || stream_req > 100) {
1242 std::stringstream ss;
1243 ss <<
"Invalid request for number of streams";
1245 logTransferEvent(
LogMask::Info, rec,
"INVALID_REQUEST", ss.str());
1246 return req.
SendSimpleResp(rec.status, NULL, NULL, generateClientErr(ss, rec).c_str(), 0);
1248 streams = stream_req == 0 ? 1 : stream_req;
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";
1260 return req.
SendSimpleResp(rec.status, NULL, NULL, generateClientErr(ss, rec).c_str(), 0);
1262 uint64_t sourceFileContentLength = 0;
1266 bool success =
false;
1267 bool mismatchDigests =
false;
1268 std::map<std::string,std::string> sourceFileReprDigest;
1269 GetRemoteFileInfoTPCPull(curl, req, sourceFileContentLength, sourceFileReprDigest, success, rec);
1273 full_url +=
"&oss.asize=" + std::to_string(sourceFileContentLength);
1274 mismatchDigests = mismatchReprDigest(sourceFileReprDigest,req,rec);
1276 if(!success || mismatchDigests) {
1283 int open_result = OpenWaitStall(*fh, full_url, mode|
SFS_O_WRONLY,
1287 int result = RedirectTransfer(curl, redirect_resource, req, fh->
error, rec);
1289 }
else if (
SFS_OK != open_result) {
1291 std::stringstream ss;
1293 if ((msg == NULL) || (*msg ==
'\0')) ss <<
"Failed to open local resource";
1298 generateClientErr(ss, rec).c_str(), 0);
1302 Stream stream(std::move(fh), streams * m_pipelining_multiplier, streams > 1 ? m_block_size : m_small_block_size, m_log);
1308 return RunCurlWithStreams(req, state, streams, rec);
1310 return RunCurlWithUpdates(curl, req, state, rec);
1318void TPCHandler::logTransferEvent(
LogMask mask,
const TPCLogRecord &rec,
1319 const std::string &event,
const std::string &message)
1321 if (!(m_log.getMsgMask() & mask)) {
return;}
1323 std::stringstream ss;
1324 ss <<
"event=" <<
event <<
", local=" << rec.local <<
", remote=" << rec.remote;
1325 if (rec.name.empty())
1326 ss <<
", user=(anonymous)";
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());
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);
1358 if (curl_global_init(CURL_GLOBAL_DEFAULT)) {
1359 log->
Emsg(
"TPCInitialize",
"libcurl failed to initialize");
1365 log->
Emsg(
"TPCInitialize",
"TPC handler requires a config filename in order to load");
1369 log->
Emsg(
"TPCInitialize",
"Will load configuration for the TPC handler from", config);
1371 }
catch (std::runtime_error &re) {
1372 log->
Emsg(
"TPCInitialize",
"Encountered a runtime failure when loading ", re.what());
XrdVERSIONINFO(XrdClGetPlugIn, XrdClGetPlugIn) extern "C"
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)
void getline(uchar *buff, int blen)
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)
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 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::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 ¶ms)
static int GetSokInfo(int fd, char *theAddr, int theALen, char &theType)
void * GetPtr(const char *varname)
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.
virtual int open(const char *fileName, XrdSfsFileOpenMode openMode, mode_t createMode, const XrdSecEntity *client=0, const char *opaque=0)=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