1087 int max_pending = 50;
1089 m_continue_queue.reset(
new HandlerQueue(max_pending));
1090 auto &queue = *m_queue.get();
1093 CURLM *multi_handle = curl_multi_init();
1094 if (multi_handle ==
nullptr) {
1095 throw std::runtime_error(
"Failed to create curl multi-handle");
1098 int running_handles = 0;
1099 time_t last_maintenance = time(NULL);
1100 CURLMcode mres = CURLM_OK;
1104 std::unordered_map<int, WaitingForBroker> broker_reqs;
1105 std::vector<struct curl_waitfd> waitfds;
1107 bool want_shutdown =
false;
1108 while (!want_shutdown) {
1109 m_last_completed_cycle.store(std::chrono::system_clock::now().time_since_epoch().count());
1110 auto oldest_op = std::chrono::system_clock::now();
1111 for (
const auto &entry : m_op_map) {
1112 OpRecord(*entry.second.first, OpKind::Update);
1113 if (entry.second.second < oldest_op) {
1114 oldest_op = entry.second.second;
1117 m_oldest_op.store(oldest_op.time_since_epoch().count());
1121 auto op = m_continue_queue->TryConsume();
1128 m_logger->
Debug(
kLogXrdClHttp,
"Ignoring continuation of operation that has already completed");
1131 m_logger->
Debug(
kLogXrdClHttp,
"Continuing the curl handle from op %p on thread %d", op.get(), getthreadid());
1132 auto curl = op->GetCurlHandle();
1133 if (!op->ContinueHandle()) {
1134 op->Fail(
XrdCl::errInternal, 0,
"Failed to continue the curl handle for the operation");
1136 op->ReleaseHandle();
1138 curl_multi_remove_handle(multi_handle, curl);
1139 curl_easy_cleanup(curl);
1140 m_op_map.erase(curl);
1142 running_handles -= 1;
1145 auto iter = m_op_map.find(curl);
1146 if (iter != m_op_map.end()) iter->second.second = std::chrono::system_clock::now();
1150 while (running_handles <
static_cast<int>(m_max_ops)) {
1151 auto op = running_handles == 0 ? queue.Consume(std::chrono::seconds(1)) : queue.TryConsume();
1155 auto curl = queue.GetHandle();
1156 if (curl ==
nullptr) {
1162 auto rv = op->Setup(curl, *
this);
1165 op->Fail(
XrdCl::errInternal, ENOMEM,
"Failed to setup the curl handle for the operation");
1168 if (!op->FinishSetup(curl)) {
1170 op->Fail(
XrdCl::errInternal, ENOMEM,
"Failed to finish setup of the curl handle for the operation");
1175 op->Fail(
XrdCl::errInternal, ENOMEM,
"Failed to setup the curl handle for the operation");
1178 op->SetContinueQueue(m_continue_queue);
1183 m_op_map[curl] = {op, std::chrono::system_clock::now()};
1187 if (op->RequiresOptions()) {
1188 std::string modified_url;
1189 std::shared_ptr<CurlOptionsOp> options_op(
1195 m_logger, op->GetConnCalloutFunc()
1201 curl = queue.GetHandle();
1202 if (curl ==
nullptr) {
1208 auto rv = options_op->Setup(curl, *
this);
1213 m_op_map[curl] = {options_op, std::chrono::system_clock::now()};
1214 OpRecord(*options_op, OpKind::Start);
1215 running_handles += 1;
1217 OpRecord(*op, OpKind::Start);
1220 auto mres = curl_multi_add_handle(multi_handle, curl);
1221 if (mres != CURLM_OK) {
1223 op->Fail(
XrdCl::errInternal, mres,
"Unable to add operation to the curl multi-handle");
1227 m_logger->
Debug(
kLogXrdClHttp,
"Added request for URL %s to worker thread for processing", op->GetUrl().c_str());
1228 running_handles += 1;
1233 time_t now = time(NULL);
1234 time_t next_maintenance = last_maintenance + m_maintenance_period.load(std::memory_order_relaxed);
1235 if (now >= next_maintenance) {
1237 m_continue_queue->Expire();
1239 getthreadid(), running_handles);
1240 last_maintenance = now;
1243 std::vector<std::pair<int, CURL *>> expired_ops;
1244 for (
const auto &entry : broker_reqs) {
1245 if (entry.second.expiry < now) {
1246 expired_ops.emplace_back(entry.first, entry.second.curl);
1249 for (
const auto &entry : expired_ops) {
1250 auto iter = m_op_map.find(entry.second);
1251 if (iter == m_op_map.end()) {
1252 m_logger->
Warning(
kLogXrdClHttp,
"Found an expired curl handle with no corresponding operation!");
1255 CurlOptionsOp *options_op =
nullptr;
1256 if ((options_op =
dynamic_cast<CurlOptionsOp*
>(iter->second.first.get())) !=
nullptr) {
1257 auto parent_op = options_op->GetOperation();
1258 bool parent_op_failed =
false;
1259 if (parent_op->IsRedirect()) {
1261 if (parent_op->Redirect(target) == CurlOperation::RedirectAction::Fail) {
1262 auto iter = m_op_map.find(options_op->GetParentCurlHandle());
1263 if (iter != m_op_map.end()) {
1266 m_op_map.erase(iter);
1267 running_handles -= 1;
1269 parent_op_failed =
true;
1271 OpRecord(*parent_op, OpKind::Start);
1274 OpRecord(*parent_op, OpKind::Start);
1276 if (!parent_op_failed){
1277 curl_multi_add_handle(multi_handle, options_op->GetParentCurlHandle());
1282 iter->second.first->ReleaseHandle();
1283 OpRecord(*(iter->second.first), OpKind::ConncallTimeout);
1284 m_op_map.erase(entry.second);
1285 curl_easy_cleanup(entry.second);
1286 running_handles -= 1;
1288 broker_reqs.erase(entry.first);
1289 m_conncall_timeout.fetch_add(1, std::memory_order_relaxed);
1297 waitfds.resize(3 + broker_reqs.size());
1299 waitfds[0].fd = queue.PollFD();
1300 waitfds[0].events = CURL_WAIT_POLLIN;
1301 waitfds[0].revents = 0;
1302 waitfds[1].fd = m_continue_queue->PollFD();
1303 waitfds[1].events = CURL_WAIT_POLLIN;
1304 waitfds[1].revents = 0;
1305 waitfds[2].fd = m_shutdown_pipe_r;
1306 waitfds[2].revents = 0;
1307 waitfds[2].events = CURL_WAIT_POLLIN | CURL_WAIT_POLLPRI;
1310 for (
const auto &entry : broker_reqs) {
1311 waitfds[idx].fd = entry.first;
1312 waitfds[idx].events = CURL_WAIT_POLLIN|CURL_WAIT_POLLPRI;
1313 waitfds[idx].revents = 0;
1318 curl_multi_timeout(multi_handle, &timeo);
1322 if (running_handles && timeo == -1) {
1327 mres = curl_multi_wait(multi_handle, &waitfds[0], waitfds.size(), 50,
nullptr);
1333 mres = curl_multi_wait(multi_handle, &waitfds[0], waitfds.size(), 50,
nullptr);
1335 if (mres != CURLM_OK) {
1340 for (
const auto &entry : waitfds) {
1342 if (waitfds[0].fd == entry.fd || waitfds[1].fd == entry.fd) {
1346 if ((waitfds[2].fd == entry.fd) && entry.revents) {
1347 want_shutdown =
true;
1350 if ((entry.revents & CURL_WAIT_POLLIN) != CURL_WAIT_POLLIN) {
1353 auto handle = broker_reqs[entry.fd].curl;
1354 auto iter = m_op_map.find(handle);
1355 if (iter == m_op_map.end()) {
1356 m_logger->
Warning(
kLogXrdClHttp,
"Internal error: broker responded on FD %d but no corresponding curl operation", entry.fd);
1357 broker_reqs.erase(entry.fd);
1358 m_conncall_errors.fetch_add(1, std::memory_order_relaxed);
1362 auto result = iter->second.first->WaitSocketCallback(err);
1366 CurlOptionsOp *options_op =
nullptr;
1367 if ((options_op =
dynamic_cast<CurlOptionsOp*
>(iter->second.first.get())) !=
nullptr) {
1368 auto parent_op = options_op->GetOperation();
1369 bool parent_op_failed =
false;
1370 if (parent_op->IsRedirect()) {
1372 if (parent_op->Redirect(target) == CurlOperation::RedirectAction::Fail) {
1373 auto iter = m_op_map.find(options_op->GetParentCurlHandle());
1374 if (iter != m_op_map.end()) {
1377 m_op_map.erase(iter);
1378 running_handles -= 1;
1380 parent_op_failed =
true;
1382 OpRecord(*parent_op, OpKind::Start);
1385 OpRecord(*parent_op, OpKind::Start);
1387 if (!parent_op_failed){
1388 curl_multi_add_handle(multi_handle, options_op->GetParentCurlHandle());
1394 m_op_map.erase(handle);
1395 broker_reqs.erase(entry.fd);
1396 m_conncall_errors.fetch_add(1, std::memory_order_relaxed);
1397 running_handles -= 1;
1399 broker_reqs.erase(entry.fd);
1400 curl_multi_add_handle(multi_handle, handle);
1401 m_conncall_success.fetch_add(1, std::memory_order_relaxed);
1407 auto mres = curl_multi_perform(multi_handle, &still_running);
1408 if (mres == CURLM_CALL_MULTI_PERFORM) {
1410 }
else if (mres != CURLM_OK) {
1418 msg = curl_multi_info_read(multi_handle, &msgq);
1419 if (msg && (msg->msg == CURLMSG_DONE)) {
1420 if (!msg->easy_handle) {
1422 mres = CURLM_BAD_EASY_HANDLE;
1425 auto iter = m_op_map.find(msg->easy_handle);
1426 if (iter == m_op_map.end()) {
1427 m_logger->
Error(
kLogXrdClHttp,
"Logic error: got a callback for an entry that doesn't exist");
1428 mres = CURLM_BAD_EASY_HANDLE;
1431 auto op = iter->second.first;
1432 auto res = msg->data.result;
1433 bool keep_handle =
false;
1434 bool waiting_on_callout =
false;
1435 if (res == CURLE_OK) {
1436 auto sc = op->GetStatusCode();
1437 OpRecord(*op, OpKind::Finish);
1440 op->Fail(httpErr.first, httpErr.second, op->GetStatusMessage());
1441 op->ReleaseHandle();
1445 CurlOptionsOp *options_op =
nullptr;
1446 if ((options_op =
dynamic_cast<CurlOptionsOp*
>(op.get())) !=
nullptr) {
1447 auto parent_op = options_op->GetOperation();
1448 bool parent_op_failed =
false;
1449 if (parent_op->IsRedirect()) {
1451 if (parent_op->Redirect(target) == CurlOperation::RedirectAction::Fail) {
1453 m_op_map.erase(options_op->GetParentCurlHandle());
1454 running_handles -= 1;
1455 parent_op_failed =
true;
1457 OpRecord(*parent_op, OpKind::Start);
1460 OpRecord(*parent_op, OpKind::Start);
1463 if (!parent_op_failed) {
1464 curl_multi_add_handle(multi_handle, options_op->GetParentCurlHandle());
1468 queue.RecycleHandle(iter->first);
1470 CurlOptionsOp *options_op =
nullptr;
1472 if ((options_op =
dynamic_cast<CurlOptionsOp*
>(op.get()))) {
1473 options_op->Success();
1474 options_op->ReleaseHandle();
1476 op = options_op->GetOperation();
1478 OpRecord(*op, OpKind::Start);
1479 curl_multi_add_handle(multi_handle, options_op->GetParentCurlHandle());
1480 curl_multi_remove_handle(multi_handle, iter->first);
1481 queue.RecycleHandle(iter->first);
1486 if (op->IsRedirect()) {
1488 switch (op->Redirect(target)) {
1489 case CurlOperation::RedirectAction::Fail:
1498 keep_handle =
false;
1500 case CurlOperation::RedirectAction::Reinvoke:
1506 OpRecord(*op, OpKind::Start);
1509 case CurlOperation::RedirectAction::ReinvokeAfterAllow:
1516 std::string modified_url;
1518 options_op =
new CurlOptionsOp(iter->first, op, target, m_logger, op->GetConnCalloutFunc());
1519 std::shared_ptr<CurlOperation> new_op(options_op);
1520 auto curl = queue.GetHandle();
1521 if (curl ==
nullptr) {
1524 keep_handle =
false;
1525 options_op =
nullptr;
1528 OpRecord(*new_op, OpKind::Start);
1530 auto rv = new_op->Setup(curl, *
this);
1533 keep_handle =
false;
1534 options_op =
nullptr;
1538 m_logger->
Debug(
kLogXrdClHttp,
"Unable to setup the curl handle for the OPTIONS operation");
1539 new_op->Fail(
XrdCl::errInternal, ENOMEM,
"Failed to setup the curl handle for the OPTIONS operation");
1541 keep_handle =
false;
1544 new_op->SetContinueQueue(m_continue_queue);
1545 m_op_map[curl] = {new_op, std::chrono::system_clock::now()};
1546 auto mres = curl_multi_add_handle(multi_handle, curl);
1547 if (mres != CURLM_OK) {
1548 m_logger->
Debug(
kLogXrdClHttp,
"Unable to add OPTIONS operation to the curl multi-handle: %s", curl_multi_strerror(mres));
1549 op->Fail(
XrdCl::errInternal, mres,
"Unable to add OPTIONS operation to the curl multi-handle");
1553 running_handles += 1;
1554 m_logger->
Debug(
kLogXrdClHttp,
"Invoking the OPTIONS operation before redirect to %s", target.c_str());
1560 int callout_socket = op->WaitSocket();
1561 if ((waiting_on_callout = callout_socket >= 0)) {
1562 auto expiry = time(
nullptr) + 20;
1563 m_logger->
Debug(
kLogXrdClHttp,
"Creating a callout wait request on socket %d", callout_socket);
1564 broker_reqs[callout_socket] = {iter->first, expiry};
1565 m_conncall_req.fetch_add(1, std::memory_order_relaxed);
1567 }
else if (options_op) {
1569 curl_multi_add_handle(multi_handle, options_op->GetParentCurlHandle());
1572 curl_multi_remove_handle(multi_handle, iter->first);
1573 if (!waiting_on_callout && !options_op) {
1574 curl_multi_add_handle(multi_handle, iter->first);
1576 }
else if (!options_op) {
1578 op->ReleaseHandle();
1580 queue.RecycleHandle(iter->first);
1583 }
else if (res == CURLE_COULDNT_CONNECT && op->UseConnectionCallout() && !op->GetTriedBoker()) {
1587 op->SetTriedBoker();
1589 int wait_socket = -1;
1590 if (!op->StartConnectionCallout(err) || (wait_socket=op->WaitSocket()) == -1) {
1591 m_logger->
Error(
kLogXrdClHttp,
"Failed to start broker-based connection: %s", err.c_str());
1592 op->ReleaseHandle();
1593 keep_handle =
false;
1595 curl_multi_remove_handle(multi_handle, iter->first);
1596 auto expiry = time(
nullptr) + 20;
1597 m_logger->
Debug(
kLogXrdClHttp,
"Curl operation requires a new TCP socket; waiting for callout to respond on socket %d", wait_socket);
1598 broker_reqs[wait_socket] = {iter->first, expiry};
1599 m_conncall_req.fetch_add(1, std::memory_order_relaxed);
1602 if (res == CURLE_ABORTED_BY_CALLBACK || res == CURLE_WRITE_ERROR) {
1605 switch (op->GetError()) {
1606 case CurlOperation::OpError::ErrHeaderTimeout:
1607 #ifdef HAVE_XPROTOCOL_TIMEREXPIRED
1614 case CurlOperation::OpError::ErrCallback: {
1615 auto [ecode,
emsg] = op->GetCallbackError();
1620 case CurlOperation::OpError::ErrOperationTimeout:
1622 OpRecord(*op, op->IsPaused() ? OpKind::ClientTimeout : OpKind::ServerTimeout);
1624 case CurlOperation::OpError::ErrTransferSlow:
1626 OpRecord(*op, OpKind::ServerTimeout);
1628 case CurlOperation::OpError::ErrTransferClientStall:
1630 OpRecord(*op, OpKind::ClientTimeout);
1632 case CurlOperation::OpError::ErrTransferStall:
1634 OpRecord(*op, OpKind::ServerTimeout);
1636 case CurlOperation::OpError::ErrNone:
1637 op->Fail(
XrdCl::errInternal, 0,
"Operation was aborted without recording an abort reason");
1641 CurlOptionsOp *options_op =
nullptr;
1642 if ((options_op =
dynamic_cast<CurlOptionsOp*
>(op.get())) !=
nullptr) {
1643 auto parent_op = options_op->GetOperation();
1644 bool parent_op_failed =
false;
1645 if (parent_op->IsRedirect()) {
1647 if (parent_op->Redirect(target) == CurlOperation::RedirectAction::Fail) {
1648 auto iter = m_op_map.find(options_op->GetParentCurlHandle());
1649 if (iter != m_op_map.end()) {
1652 m_op_map.erase(iter);
1653 running_handles -= 1;
1655 parent_op_failed =
true;
1657 OpRecord(*parent_op, OpKind::Start);
1660 OpRecord(*parent_op, OpKind::Start);
1662 if (!parent_op_failed){
1663 curl_multi_add_handle(multi_handle, options_op->GetParentCurlHandle());
1668 const auto curl_err = op->GetCurlErrorMessage();
1669 const char *curl_easy_err = curl_easy_strerror(res);
1670 const std::string fail_err = !curl_err.empty() ? curl_err : curl_easy_err;
1671 m_logger->
Debug(
kLogXrdClHttp,
"Curl generated an error: %s (%d)", fail_err.c_str(), res);
1672 op->Fail(xrdCode.first, xrdCode.second, fail_err);
1674 CurlOptionsOp *options_op =
nullptr;
1675 if ((options_op =
dynamic_cast<CurlOptionsOp*
>(op.get())) !=
nullptr) {
1676 auto parent_op = options_op->GetOperation();
1677 bool parent_op_failed =
false;
1678 if (parent_op->IsRedirect()) {
1680 if (parent_op->Redirect(target) == CurlOperation::RedirectAction::Fail) {
1681 auto iter = m_op_map.find(options_op->GetParentCurlHandle());
1682 if (iter != m_op_map.end()) {
1685 m_op_map.erase(iter);
1686 running_handles -= 1;
1688 parent_op_failed =
true;
1691 if (!parent_op_failed){
1692 curl_multi_add_handle(multi_handle, options_op->GetParentCurlHandle());
1696 op->ReleaseHandle();
1699 curl_multi_remove_handle(multi_handle, iter->first);
1700 if (res != CURLE_OK) {
1701 curl_easy_cleanup(iter->first);
1703 for (
auto &req : broker_reqs) {
1704 if (req.second.curl == iter->first) {
1705 m_logger->
Warning(
kLogXrdClHttp,
"Curl handle finished while a broker operation was outstanding");
1706 m_conncall_errors.fetch_add(1, std::memory_order_relaxed);
1709 m_op_map.erase(iter);
1710 running_handles -= 1;
1716 for (
auto map_entry : m_op_map) {
1721 if (multi_handle && map_entry.first) curl_multi_remove_handle(multi_handle, map_entry.first);
1724 m_queue->ReleaseHandles();
1725 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