1091 int max_pending = 50;
1093 m_continue_queue.reset(
new HandlerQueue(max_pending));
1094 auto &queue = *m_queue.get();
1097 CURLM *multi_handle = curl_multi_init();
1098 if (multi_handle ==
nullptr) {
1099 throw std::runtime_error(
"Failed to create curl multi-handle");
1102 int running_handles = 0;
1103 time_t last_maintenance = time(NULL);
1104 CURLMcode mres = CURLM_OK;
1108 std::unordered_map<int, WaitingForBroker> broker_reqs;
1109 std::vector<struct curl_waitfd> waitfds;
1111 bool want_shutdown =
false;
1112 while (!want_shutdown) {
1113 m_last_completed_cycle.store(std::chrono::system_clock::now().time_since_epoch().count());
1114 auto oldest_op = std::chrono::system_clock::now();
1115 for (
const auto &entry : m_op_map) {
1116 OpRecord(*entry.second.first, OpKind::Update);
1117 if (entry.second.second < oldest_op) {
1118 oldest_op = entry.second.second;
1121 m_oldest_op.store(oldest_op.time_since_epoch().count());
1125 auto op = m_continue_queue->TryConsume();
1132 m_logger->
Debug(
kLogXrdClHttp,
"Ignoring continuation of operation that has already completed");
1135 m_logger->
Debug(
kLogXrdClHttp,
"Continuing the curl handle from op %p on thread %d", op.get(), getthreadid());
1136 auto curl = op->GetCurlHandle();
1137 if (!op->ContinueHandle()) {
1138 op->Fail(
XrdCl::errInternal, 0,
"Failed to continue the curl handle for the operation");
1140 op->ReleaseHandle();
1142 curl_multi_remove_handle(multi_handle, curl);
1143 curl_easy_cleanup(curl);
1144 m_op_map.erase(curl);
1146 running_handles -= 1;
1149 auto iter = m_op_map.find(curl);
1150 if (iter != m_op_map.end()) iter->second.second = std::chrono::system_clock::now();
1154 while (running_handles <
static_cast<int>(m_max_ops)) {
1155 auto op = running_handles == 0 ? queue.Consume(std::chrono::seconds(1)) : queue.TryConsume();
1159 auto curl = queue.GetHandle();
1160 if (curl ==
nullptr) {
1166 auto rv = op->Setup(curl, *
this);
1169 op->Fail(
XrdCl::errInternal, ENOMEM,
"Failed to setup the curl handle for the operation");
1172 if (!op->FinishSetup(curl)) {
1174 op->Fail(
XrdCl::errInternal, ENOMEM,
"Failed to finish setup of the curl handle for the operation");
1179 op->Fail(
XrdCl::errInternal, ENOMEM,
"Failed to setup the curl handle for the operation");
1182 op->SetContinueQueue(m_continue_queue);
1187 m_op_map[curl] = {op, std::chrono::system_clock::now()};
1191 if (op->RequiresOptions()) {
1192 std::string modified_url;
1193 std::shared_ptr<CurlOptionsOp> options_op(
1199 m_logger, op->GetConnCalloutFunc()
1205 curl = queue.GetHandle();
1206 if (curl ==
nullptr) {
1212 auto rv = options_op->Setup(curl, *
this);
1217 m_op_map[curl] = {options_op, std::chrono::system_clock::now()};
1218 OpRecord(*options_op, OpKind::Start);
1219 running_handles += 1;
1221 OpRecord(*op, OpKind::Start);
1224 auto mres = curl_multi_add_handle(multi_handle, curl);
1225 if (mres != CURLM_OK) {
1227 op->Fail(
XrdCl::errInternal, mres,
"Unable to add operation to the curl multi-handle");
1231 m_logger->
Debug(
kLogXrdClHttp,
"Added request for URL %s to worker thread for processing", op->GetUrl().c_str());
1232 running_handles += 1;
1237 time_t now = time(NULL);
1238 time_t next_maintenance = last_maintenance + m_maintenance_period.load(std::memory_order_relaxed);
1239 if (now >= next_maintenance) {
1241 m_continue_queue->Expire();
1243 getthreadid(), running_handles);
1244 last_maintenance = now;
1247 std::vector<std::pair<int, CURL *>> expired_ops;
1248 for (
const auto &entry : broker_reqs) {
1249 if (entry.second.expiry < now) {
1250 expired_ops.emplace_back(entry.first, entry.second.curl);
1253 for (
const auto &entry : expired_ops) {
1254 auto iter = m_op_map.find(entry.second);
1255 if (iter == m_op_map.end()) {
1256 m_logger->
Warning(
kLogXrdClHttp,
"Found an expired curl handle with no corresponding operation!");
1259 CurlOptionsOp *options_op =
nullptr;
1260 if ((options_op =
dynamic_cast<CurlOptionsOp*
>(iter->second.first.get())) !=
nullptr) {
1261 auto parent_op = options_op->GetOperation();
1262 bool parent_op_failed =
false;
1263 if (parent_op->IsRedirect()) {
1265 if (parent_op->Redirect(target) == CurlOperation::RedirectAction::Fail) {
1266 auto iter = m_op_map.find(options_op->GetParentCurlHandle());
1267 if (iter != m_op_map.end()) {
1270 m_op_map.erase(iter);
1271 running_handles -= 1;
1273 parent_op_failed =
true;
1275 OpRecord(*parent_op, OpKind::Start);
1278 OpRecord(*parent_op, OpKind::Start);
1280 if (!parent_op_failed){
1281 curl_multi_add_handle(multi_handle, options_op->GetParentCurlHandle());
1286 iter->second.first->ReleaseHandle();
1287 OpRecord(*(iter->second.first), OpKind::ConncallTimeout);
1288 m_op_map.erase(entry.second);
1289 curl_easy_cleanup(entry.second);
1290 running_handles -= 1;
1292 broker_reqs.erase(entry.first);
1293 m_conncall_timeout.fetch_add(1, std::memory_order_relaxed);
1301 waitfds.resize(3 + broker_reqs.size());
1303 waitfds[0].fd = queue.PollFD();
1304 waitfds[0].events = CURL_WAIT_POLLIN;
1305 waitfds[0].revents = 0;
1306 waitfds[1].fd = m_continue_queue->PollFD();
1307 waitfds[1].events = CURL_WAIT_POLLIN;
1308 waitfds[1].revents = 0;
1309 waitfds[2].fd = m_shutdown_pipe_r;
1310 waitfds[2].revents = 0;
1311 waitfds[2].events = CURL_WAIT_POLLIN | CURL_WAIT_POLLPRI;
1314 for (
const auto &entry : broker_reqs) {
1315 waitfds[idx].fd = entry.first;
1316 waitfds[idx].events = CURL_WAIT_POLLIN|CURL_WAIT_POLLPRI;
1317 waitfds[idx].revents = 0;
1322 curl_multi_timeout(multi_handle, &timeo);
1326 if (running_handles && timeo == -1) {
1331 mres = curl_multi_wait(multi_handle, &waitfds[0], waitfds.size(), 50,
nullptr);
1337 mres = curl_multi_wait(multi_handle, &waitfds[0], waitfds.size(), 50,
nullptr);
1339 if (mres != CURLM_OK) {
1344 for (
const auto &entry : waitfds) {
1346 if (waitfds[0].fd == entry.fd || waitfds[1].fd == entry.fd) {
1350 if ((waitfds[2].fd == entry.fd) && entry.revents) {
1351 want_shutdown =
true;
1354 if ((entry.revents & CURL_WAIT_POLLIN) != CURL_WAIT_POLLIN) {
1357 auto handle = broker_reqs[entry.fd].curl;
1358 auto iter = m_op_map.find(handle);
1359 if (iter == m_op_map.end()) {
1360 m_logger->
Warning(
kLogXrdClHttp,
"Internal error: broker responded on FD %d but no corresponding curl operation", entry.fd);
1361 broker_reqs.erase(entry.fd);
1362 m_conncall_errors.fetch_add(1, std::memory_order_relaxed);
1366 auto result = iter->second.first->WaitSocketCallback(err);
1370 CurlOptionsOp *options_op =
nullptr;
1371 if ((options_op =
dynamic_cast<CurlOptionsOp*
>(iter->second.first.get())) !=
nullptr) {
1372 auto parent_op = options_op->GetOperation();
1373 bool parent_op_failed =
false;
1374 if (parent_op->IsRedirect()) {
1376 if (parent_op->Redirect(target) == CurlOperation::RedirectAction::Fail) {
1377 auto iter = m_op_map.find(options_op->GetParentCurlHandle());
1378 if (iter != m_op_map.end()) {
1381 m_op_map.erase(iter);
1382 running_handles -= 1;
1384 parent_op_failed =
true;
1386 OpRecord(*parent_op, OpKind::Start);
1389 OpRecord(*parent_op, OpKind::Start);
1391 if (!parent_op_failed){
1392 curl_multi_add_handle(multi_handle, options_op->GetParentCurlHandle());
1398 m_op_map.erase(handle);
1399 broker_reqs.erase(entry.fd);
1400 m_conncall_errors.fetch_add(1, std::memory_order_relaxed);
1401 running_handles -= 1;
1403 broker_reqs.erase(entry.fd);
1404 curl_multi_add_handle(multi_handle, handle);
1405 m_conncall_success.fetch_add(1, std::memory_order_relaxed);
1411 auto mres = curl_multi_perform(multi_handle, &still_running);
1412 if (mres == CURLM_CALL_MULTI_PERFORM) {
1414 }
else if (mres != CURLM_OK) {
1422 msg = curl_multi_info_read(multi_handle, &msgq);
1423 if (msg && (msg->msg == CURLMSG_DONE)) {
1424 if (!msg->easy_handle) {
1426 mres = CURLM_BAD_EASY_HANDLE;
1429 auto iter = m_op_map.find(msg->easy_handle);
1430 if (iter == m_op_map.end()) {
1431 m_logger->
Error(
kLogXrdClHttp,
"Logic error: got a callback for an entry that doesn't exist");
1432 mres = CURLM_BAD_EASY_HANDLE;
1435 auto op = iter->second.first;
1436 auto res = msg->data.result;
1437 bool keep_handle =
false;
1438 bool waiting_on_callout =
false;
1439 if (res == CURLE_OK) {
1440 auto sc = op->GetStatusCode();
1441 OpRecord(*op, OpKind::Finish);
1444 op->Fail(httpErr.first, httpErr.second, op->GetStatusMessage());
1445 op->ReleaseHandle();
1449 CurlOptionsOp *options_op =
nullptr;
1450 if ((options_op =
dynamic_cast<CurlOptionsOp*
>(op.get())) !=
nullptr) {
1451 auto parent_op = options_op->GetOperation();
1452 bool parent_op_failed =
false;
1453 if (parent_op->IsRedirect()) {
1455 if (parent_op->Redirect(target) == CurlOperation::RedirectAction::Fail) {
1457 m_op_map.erase(options_op->GetParentCurlHandle());
1458 running_handles -= 1;
1459 parent_op_failed =
true;
1461 OpRecord(*parent_op, OpKind::Start);
1464 OpRecord(*parent_op, OpKind::Start);
1467 if (!parent_op_failed) {
1468 curl_multi_add_handle(multi_handle, options_op->GetParentCurlHandle());
1472 queue.RecycleHandle(iter->first);
1474 CurlOptionsOp *options_op =
nullptr;
1476 if ((options_op =
dynamic_cast<CurlOptionsOp*
>(op.get()))) {
1477 options_op->Success();
1478 options_op->ReleaseHandle();
1480 op = options_op->GetOperation();
1482 OpRecord(*op, OpKind::Start);
1483 curl_multi_add_handle(multi_handle, options_op->GetParentCurlHandle());
1484 curl_multi_remove_handle(multi_handle, iter->first);
1485 queue.RecycleHandle(iter->first);
1490 if (op->IsRedirect()) {
1492 switch (op->Redirect(target)) {
1493 case CurlOperation::RedirectAction::Fail:
1502 keep_handle =
false;
1504 case CurlOperation::RedirectAction::Reinvoke:
1510 OpRecord(*op, OpKind::Start);
1513 case CurlOperation::RedirectAction::ReinvokeAfterAllow:
1520 std::string modified_url;
1522 options_op =
new CurlOptionsOp(iter->first, op, target, m_logger, op->GetConnCalloutFunc());
1523 std::shared_ptr<CurlOperation> new_op(options_op);
1524 auto curl = queue.GetHandle();
1525 if (curl ==
nullptr) {
1528 keep_handle =
false;
1529 options_op =
nullptr;
1532 OpRecord(*new_op, OpKind::Start);
1534 auto rv = new_op->Setup(curl, *
this);
1537 keep_handle =
false;
1538 options_op =
nullptr;
1542 m_logger->
Debug(
kLogXrdClHttp,
"Unable to setup the curl handle for the OPTIONS operation");
1543 new_op->Fail(
XrdCl::errInternal, ENOMEM,
"Failed to setup the curl handle for the OPTIONS operation");
1545 keep_handle =
false;
1548 new_op->SetContinueQueue(m_continue_queue);
1549 m_op_map[curl] = {new_op, std::chrono::system_clock::now()};
1550 auto mres = curl_multi_add_handle(multi_handle, curl);
1551 if (mres != CURLM_OK) {
1552 m_logger->
Debug(
kLogXrdClHttp,
"Unable to add OPTIONS operation to the curl multi-handle: %s", curl_multi_strerror(mres));
1553 op->Fail(
XrdCl::errInternal, mres,
"Unable to add OPTIONS operation to the curl multi-handle");
1557 running_handles += 1;
1558 m_logger->
Debug(
kLogXrdClHttp,
"Invoking the OPTIONS operation before redirect to %s", target.c_str());
1564 int callout_socket = op->WaitSocket();
1565 if ((waiting_on_callout = callout_socket >= 0)) {
1566 auto expiry = time(
nullptr) + 20;
1567 m_logger->
Debug(
kLogXrdClHttp,
"Creating a callout wait request on socket %d", callout_socket);
1568 broker_reqs[callout_socket] = {iter->first, expiry};
1569 m_conncall_req.fetch_add(1, std::memory_order_relaxed);
1571 }
else if (options_op) {
1573 curl_multi_add_handle(multi_handle, options_op->GetParentCurlHandle());
1576 curl_multi_remove_handle(multi_handle, iter->first);
1577 if (!waiting_on_callout && !options_op) {
1578 curl_multi_add_handle(multi_handle, iter->first);
1580 }
else if (!options_op) {
1582 op->ReleaseHandle();
1584 queue.RecycleHandle(iter->first);
1587 }
else if (res == CURLE_COULDNT_CONNECT && op->UseConnectionCallout() && !op->GetTriedBoker()) {
1591 op->SetTriedBoker();
1593 int wait_socket = -1;
1594 if (!op->StartConnectionCallout(err) || (wait_socket=op->WaitSocket()) == -1) {
1595 m_logger->
Error(
kLogXrdClHttp,
"Failed to start broker-based connection: %s", err.c_str());
1596 op->ReleaseHandle();
1597 keep_handle =
false;
1599 curl_multi_remove_handle(multi_handle, iter->first);
1600 auto expiry = time(
nullptr) + 20;
1601 m_logger->
Debug(
kLogXrdClHttp,
"Curl operation requires a new TCP socket; waiting for callout to respond on socket %d", wait_socket);
1602 broker_reqs[wait_socket] = {iter->first, expiry};
1603 m_conncall_req.fetch_add(1, std::memory_order_relaxed);
1606 if (res == CURLE_ABORTED_BY_CALLBACK || res == CURLE_WRITE_ERROR) {
1609 switch (op->GetError()) {
1610 case CurlOperation::OpError::ErrHeaderTimeout:
1611 #ifdef HAVE_XPROTOCOL_TIMEREXPIRED
1618 case CurlOperation::OpError::ErrCallback: {
1619 auto [ecode,
emsg] = op->GetCallbackError();
1624 case CurlOperation::OpError::ErrOperationTimeout:
1626 OpRecord(*op, op->IsPaused() ? OpKind::ClientTimeout : OpKind::ServerTimeout);
1628 case CurlOperation::OpError::ErrTransferSlow:
1630 OpRecord(*op, OpKind::ServerTimeout);
1632 case CurlOperation::OpError::ErrTransferClientStall:
1634 OpRecord(*op, OpKind::ClientTimeout);
1636 case CurlOperation::OpError::ErrTransferStall:
1638 OpRecord(*op, OpKind::ServerTimeout);
1640 case CurlOperation::OpError::ErrNone:
1641 op->Fail(
XrdCl::errInternal, 0,
"Operation was aborted without recording an abort reason");
1645 CurlOptionsOp *options_op =
nullptr;
1646 if ((options_op =
dynamic_cast<CurlOptionsOp*
>(op.get())) !=
nullptr) {
1647 auto parent_op = options_op->GetOperation();
1648 bool parent_op_failed =
false;
1649 if (parent_op->IsRedirect()) {
1651 if (parent_op->Redirect(target) == CurlOperation::RedirectAction::Fail) {
1652 auto iter = m_op_map.find(options_op->GetParentCurlHandle());
1653 if (iter != m_op_map.end()) {
1656 m_op_map.erase(iter);
1657 running_handles -= 1;
1659 parent_op_failed =
true;
1661 OpRecord(*parent_op, OpKind::Start);
1664 OpRecord(*parent_op, OpKind::Start);
1666 if (!parent_op_failed){
1667 curl_multi_add_handle(multi_handle, options_op->GetParentCurlHandle());
1672 const auto curl_err = op->GetCurlErrorMessage();
1673 const char *curl_easy_err = curl_easy_strerror(res);
1674 const std::string fail_err = !curl_err.empty() ? curl_err : curl_easy_err;
1675 m_logger->
Debug(
kLogXrdClHttp,
"Curl generated an error: %s (%d)", fail_err.c_str(), res);
1676 op->Fail(xrdCode.first, xrdCode.second, fail_err);
1678 CurlOptionsOp *options_op =
nullptr;
1679 if ((options_op =
dynamic_cast<CurlOptionsOp*
>(op.get())) !=
nullptr) {
1680 auto parent_op = options_op->GetOperation();
1681 bool parent_op_failed =
false;
1682 if (parent_op->IsRedirect()) {
1684 if (parent_op->Redirect(target) == CurlOperation::RedirectAction::Fail) {
1685 auto iter = m_op_map.find(options_op->GetParentCurlHandle());
1686 if (iter != m_op_map.end()) {
1689 m_op_map.erase(iter);
1690 running_handles -= 1;
1692 parent_op_failed =
true;
1695 if (!parent_op_failed){
1696 curl_multi_add_handle(multi_handle, options_op->GetParentCurlHandle());
1700 op->ReleaseHandle();
1703 curl_multi_remove_handle(multi_handle, iter->first);
1704 if (res != CURLE_OK) {
1705 curl_easy_cleanup(iter->first);
1707 for (
auto &req : broker_reqs) {
1708 if (req.second.curl == iter->first) {
1709 m_logger->
Warning(
kLogXrdClHttp,
"Curl handle finished while a broker operation was outstanding");
1710 m_conncall_errors.fetch_add(1, std::memory_order_relaxed);
1713 m_op_map.erase(iter);
1714 running_handles -= 1;
1720 for (
auto map_entry : m_op_map) {
1725 if (multi_handle && map_entry.first) curl_multi_remove_handle(multi_handle, map_entry.first);
1728 m_queue->ReleaseHandles();
1729 curl_multi_cleanup(multi_handle);
std::pair< uint16_t, uint32_t > CurlCodeConvert(CURLcode res)
int emsg(int rc, char *msg)
static void CleanupDnsCache()
static std::string_view GetUrlKey(const std::string &url, std::string &modified_url)
bool GetInt(const std::string &key, int &value)
void Error(uint64_t topic, const char *format,...)
Report an error.
void Warning(uint64_t topic, const char *format,...)
Report a warning.
void Debug(uint64_t topic, const char *format,...)
Print a debug message.
std::pair< uint16_t, uint32_t > HTTPStatusConvert(unsigned status)
bool HTTPStatusIsError(unsigned status)
const uint16_t errErrorResponse
const uint16_t errOperationExpired
const uint16_t errInternal
Internal error.
const uint16_t errConnectionError
const uint64_t kLogXrdClHttp