1226 {
1227 int max_pending = 50;
1230 auto &queue = *m_queue.get();
1232
1233 CURLM *multi_handle = curl_multi_init();
1234 if (multi_handle == nullptr) {
1235 throw std::runtime_error("Failed to create curl multi-handle");
1236 }
1237
1238 int running_handles = 0;
1239 time_t last_maintenance = time(NULL);
1240 CURLMcode mres = CURLM_OK;
1241
1242
1243
1244 std::unordered_map<int, WaitingForBroker> broker_reqs;
1245 std::vector<struct curl_waitfd> waitfds;
1246
1247 bool want_shutdown = false;
1248 while (!want_shutdown) {
1249 m_last_completed_cycle.store(std::chrono::system_clock::now().time_since_epoch().count());
1250 auto oldest_op = std::chrono::system_clock::now();
1251 for (const auto &entry : m_op_map) {
1252 OpRecord(*entry.second.first, OpKind::Update);
1253 if (entry.second.second < oldest_op) {
1254 oldest_op = entry.second.second;
1255 }
1256 }
1257 m_oldest_op.store(oldest_op.time_since_epoch().count());
1258
1259
1260 while (true) {
1261 auto op = m_continue_queue->TryConsume();
1262 if (!op) {
1263 break;
1264 }
1265
1266
1267 if (op->IsDone()) {
1268 m_logger->Debug(
kLogXrdClHttp,
"Ignoring continuation of operation that has already completed");
1269 continue;
1270 }
1271 m_logger->Debug(
kLogXrdClHttp,
"Continuing the curl handle from op %p on thread %d", op.get(), getthreadid());
1272 auto curl = op->GetCurlHandle();
1273 if (!op->ContinueHandle()) {
1274 op->Fail(
XrdCl::errInternal, 0,
"Failed to continue the curl handle for the operation");
1275 OpRecord(*op, OpKind::Error);
1276 op->ReleaseHandle();
1277 if (curl) {
1278 curl_multi_remove_handle(multi_handle, curl);
1279 curl_easy_cleanup(curl);
1280 m_op_map.erase(curl);
1281 }
1282 running_handles -= 1;
1283 continue;
1284 } else {
1285 auto iter = m_op_map.find(curl);
1286 if (iter != m_op_map.end()) iter->second.second = std::chrono::system_clock::now();
1287 }
1288 }
1289
1290 while (running_handles < static_cast<int>(m_max_ops)) {
1291 auto op = running_handles == 0 ? queue.Consume(std::chrono::seconds(1)) : queue.TryConsume();
1292 if (!op) {
1293 break;
1294 }
1295 auto curl = queue.GetHandle();
1296 if (curl == nullptr) {
1297 m_logger->Debug(
kLogXrdClHttp,
"Unable to allocate a curl handle");
1299 continue;
1300 }
1301 try {
1302 auto rv = op->Setup(curl, *this);
1303 if (!rv) {
1304 m_logger->Debug(
kLogXrdClHttp,
"Failed to setup the curl handle");
1305 op->Fail(
XrdCl::errInternal, ENOMEM,
"Failed to setup the curl handle for the operation");
1306 continue;
1307 }
1308 if (!op->FinishSetup(curl)) {
1309 m_logger->Debug(
kLogXrdClHttp,
"Failed to finish setup of the curl handle");
1310 op->Fail(
XrdCl::errInternal, ENOMEM,
"Failed to finish setup of the curl handle for the operation");
1311 continue;
1312 }
1313 } catch (...) {
1314 m_logger->Debug(
kLogXrdClHttp,
"Unable to setup the curl handle");
1315 op->Fail(
XrdCl::errInternal, ENOMEM,
"Failed to setup the curl handle for the operation");
1316 continue;
1317 }
1318 op->SetContinueQueue(m_continue_queue);
1319
1320 if (op->IsDone()) {
1321 op->ReleaseHandle();
1322 queue.RecycleHandle(curl);
1323 continue;
1324 }
1325 m_op_map[curl] = {op, std::chrono::system_clock::now()};
1326
1327
1328
1329 if (op->RequiresOptions()) {
1330 std::string modified_url;
1331 std::shared_ptr<CurlOptionsOp> options_op(
1333 curl, op,
1334 std::string(
1336 ),
1337 m_logger, op->GetConnCalloutFunc()
1338 )
1339 );
1340
1341
1342
1343 curl = queue.GetHandle();
1344 if (curl == nullptr) {
1345 m_logger->Debug(
kLogXrdClHttp,
"Unable to allocate a curl handle");
1347 OpRecord(*op, OpKind::Error);
1348 continue;
1349 }
1350 auto rv = options_op->Setup(curl, *this);
1351 if (!rv) {
1352 m_logger->Debug(
kLogXrdClHttp,
"Failed to allocate a curl handle for OPTIONS");
1353 continue;
1354 }
1355 m_op_map[curl] = {options_op, std::chrono::system_clock::now()};
1356 OpRecord(*options_op, OpKind::Start);
1357 running_handles += 1;
1358 } else {
1359 OpRecord(*op, OpKind::Start);
1360 }
1361
1362 auto mres = curl_multi_add_handle(multi_handle, curl);
1363 if (mres != CURLM_OK) {
1364 m_logger->Debug(
kLogXrdClHttp,
"Unable to add operation to the curl multi-handle");
1365 op->Fail(
XrdCl::errInternal, mres,
"Unable to add operation to the curl multi-handle");
1366 OpRecord(*op, OpKind::Error);
1367 continue;
1368 }
1369 m_logger->Debug(
kLogXrdClHttp,
"Added request for URL %s to worker thread for processing", op->GetUrl().c_str());
1370 running_handles += 1;
1371 }
1372
1373
1374
1375 time_t now = time(NULL);
1376 time_t next_maintenance = last_maintenance + m_maintenance_period.load(std::memory_order_relaxed);
1377 if (now >= next_maintenance) {
1378 m_queue->Expire();
1379 m_continue_queue->Expire();
1380 m_logger->Debug(
kLogXrdClHttp,
"Curl worker thread %d is running %d operations",
1381 getthreadid(), running_handles);
1382 last_maintenance = now;
1383
1384
1385 std::vector<std::pair<int, CURL *>> expired_ops;
1386 for (const auto &entry : broker_reqs) {
1387 if (entry.second.expiry < now) {
1388 expired_ops.emplace_back(entry.first, entry.second.curl);
1389 }
1390 }
1391 for (const auto &entry : expired_ops) {
1392 auto iter = m_op_map.find(entry.second);
1393 if (iter == m_op_map.end()) {
1394 m_logger->Warning(
kLogXrdClHttp,
"Found an expired curl handle with no corresponding operation!");
1395 } else {
1396
1397 CurlOptionsOp *options_op = nullptr;
1398 if ((options_op = dynamic_cast<CurlOptionsOp*>(iter->second.first.get())) != nullptr) {
1400 bool parent_op_failed = false;
1401 if (parent_op->IsRedirect()) {
1402 std::string target;
1405 if (iter != m_op_map.end()) {
1406 OpRecord(*iter->second.first, OpKind::Error);
1408 m_op_map.erase(iter);
1409 running_handles -= 1;
1410 }
1411 parent_op_failed = true;
1412 } else {
1413 OpRecord(*parent_op, OpKind::Start);
1414 }
1415 } else {
1416 OpRecord(*parent_op, OpKind::Start);
1417 }
1418 if (!parent_op_failed){
1420 }
1421 }
1422
1424 iter->second.first->ReleaseHandle();
1425 OpRecord(*(iter->second.first), OpKind::ConncallTimeout);
1426 m_op_map.erase(entry.second);
1427 curl_easy_cleanup(entry.second);
1428 running_handles -= 1;
1429 }
1430 broker_reqs.erase(entry.first);
1431 m_conncall_timeout.fetch_add(1, std::memory_order_relaxed);
1432 }
1433
1434
1436 }
1437
1438 waitfds.clear();
1439 waitfds.resize(3 + broker_reqs.size());
1440
1441 waitfds[0].fd = queue.PollFD();
1442 waitfds[0].events = CURL_WAIT_POLLIN;
1443 waitfds[0].revents = 0;
1444 waitfds[1].fd = m_continue_queue->PollFD();
1445 waitfds[1].events = CURL_WAIT_POLLIN;
1446 waitfds[1].revents = 0;
1447 waitfds[2].fd = m_shutdown_pipe_r;
1448 waitfds[2].revents = 0;
1449 waitfds[2].events = CURL_WAIT_POLLIN | CURL_WAIT_POLLPRI;
1450
1451 int idx = 3;
1452 for (const auto &entry : broker_reqs) {
1453 waitfds[idx].fd = entry.first;
1454 waitfds[idx].events = CURL_WAIT_POLLIN|CURL_WAIT_POLLPRI;
1455 waitfds[idx].revents = 0;
1456 idx += 1;
1457 }
1458
1459 long timeo;
1460 curl_multi_timeout(multi_handle, &timeo);
1461
1462
1463
1464 if (running_handles && timeo == -1) {
1465
1466
1467
1468
1469 mres = curl_multi_wait(multi_handle, &waitfds[0], waitfds.size(), 50, nullptr);
1470 } else {
1471
1472
1473
1474
1475 mres = curl_multi_wait(multi_handle, &waitfds[0], waitfds.size(), 50, nullptr);
1476 }
1477 if (mres != CURLM_OK) {
1478 m_logger->Warning(
kLogXrdClHttp,
"Failed to wait on multi-handle: %d", mres);
1479 }
1480
1481
1482 for (const auto &entry : waitfds) {
1483
1484 if (waitfds[0].fd == entry.fd || waitfds[1].fd == entry.fd) {
1485 continue;
1486 }
1487
1488 if ((waitfds[2].fd == entry.fd) && entry.revents) {
1489 want_shutdown = true;
1490 break;
1491 }
1492 if ((entry.revents & CURL_WAIT_POLLIN) != CURL_WAIT_POLLIN) {
1493 continue;
1494 }
1495 auto handle = broker_reqs[entry.fd].curl;
1496 auto iter = m_op_map.find(handle);
1497 if (iter == m_op_map.end()) {
1498 m_logger->Warning(
kLogXrdClHttp,
"Internal error: broker responded on FD %d but no corresponding curl operation", entry.fd);
1499 broker_reqs.erase(entry.fd);
1500 m_conncall_errors.fetch_add(1, std::memory_order_relaxed);
1501 continue;
1502 }
1503 std::string err;
1504 auto result = iter->second.first->WaitSocketCallback(err);
1505 if (result == -1) {
1506 m_logger->Warning(
kLogXrdClHttp,
"Error when invoking the broker callback: %s", err.c_str());
1507
1508 CurlOptionsOp *options_op = nullptr;
1509 if ((options_op = dynamic_cast<CurlOptionsOp*>(iter->second.first.get())) != nullptr) {
1511 bool parent_op_failed = false;
1512 if (parent_op->IsRedirect()) {
1513 std::string target;
1516 if (iter != m_op_map.end()) {
1517 OpRecord(*iter->second.first, OpKind::Error);
1519 m_op_map.erase(iter);
1520 running_handles -= 1;
1521 }
1522 parent_op_failed = true;
1523 } else {
1524 OpRecord(*parent_op, OpKind::Start);
1525 }
1526 } else {
1527 OpRecord(*parent_op, OpKind::Start);
1528 }
1529 if (!parent_op_failed){
1531 }
1532 }
1533
1535 OpRecord(*iter->second.first, OpKind::Error);
1536 m_op_map.erase(handle);
1537 broker_reqs.erase(entry.fd);
1538 m_conncall_errors.fetch_add(1, std::memory_order_relaxed);
1539 running_handles -= 1;
1540 } else {
1541 broker_reqs.erase(entry.fd);
1542 curl_multi_add_handle(multi_handle, handle);
1543 m_conncall_success.fetch_add(1, std::memory_order_relaxed);
1544 }
1545 }
1546
1547
1548 int still_running;
1549 auto mres = curl_multi_perform(multi_handle, &still_running);
1550 if (mres == CURLM_CALL_MULTI_PERFORM) {
1551 continue;
1552 } else if (mres != CURLM_OK) {
1553 m_logger->Warning(
kLogXrdClHttp,
"Failed to perform multi-handle operation: %d", mres);
1554 break;
1555 }
1556
1557 CURLMsg *msg;
1558 do {
1559 int msgq = 0;
1560 msg = curl_multi_info_read(multi_handle, &msgq);
1561 if (msg && (msg->msg == CURLMSG_DONE)) {
1562 if (!msg->easy_handle) {
1563 m_logger->Warning(
kLogXrdClHttp,
"Logic error: got a callback for a null handle");
1564 mres = CURLM_BAD_EASY_HANDLE;
1565 break;
1566 }
1567 auto iter = m_op_map.find(msg->easy_handle);
1568 if (iter == m_op_map.end()) {
1569 m_logger->Error(
kLogXrdClHttp,
"Logic error: got a callback for an entry that doesn't exist");
1570 mres = CURLM_BAD_EASY_HANDLE;
1571 break;
1572 }
1573 auto op = iter->second.first;
1574 auto res = msg->data.result;
1575 bool keep_handle = false;
1576 bool waiting_on_callout = false;
1577 if (res == CURLE_OK) {
1578 auto sc = op->GetStatusCode();
1579 OpRecord(*op, OpKind::Finish);
1582 op->Fail(httpErr.first, httpErr.second, op->GetStatusMessage());
1583 op->ReleaseHandle();
1584
1585
1586
1587 CurlOptionsOp *options_op = nullptr;
1588 if ((options_op = dynamic_cast<CurlOptionsOp*>(op.get())) != nullptr) {
1590 bool parent_op_failed = false;
1591 if (parent_op->IsRedirect()) {
1592 std::string target;
1594 OpRecord(*parent_op, OpKind::Error);
1596 running_handles -= 1;
1597 parent_op_failed = true;
1598 } else {
1599 OpRecord(*parent_op, OpKind::Start);
1600 }
1601 } else {
1602 OpRecord(*parent_op, OpKind::Start);
1603 }
1604
1605 if (!parent_op_failed) {
1607 }
1608 }
1609
1610 queue.RecycleHandle(iter->first);
1611 } else {
1612 CurlOptionsOp *options_op = nullptr;
1613
1614 if ((options_op = dynamic_cast<CurlOptionsOp*>(op.get()))) {
1617
1619 op->OptionsDone();
1620 OpRecord(*op, OpKind::Start);
1622 curl_multi_remove_handle(multi_handle, iter->first);
1623 queue.RecycleHandle(iter->first);
1624 }
1625
1626
1627
1628 if (op->IsRedirect()) {
1629 std::string target;
1630 switch (op->Redirect(target)) {
1632 if (options_op) {
1633
1634
1635
1636
1637
1638 OpRecord(*op, OpKind::Error);
1639 }
1640 keep_handle = false;
1641 break;
1643 if (!options_op) {
1644
1645
1646
1647 keep_handle = true;
1648 OpRecord(*op, OpKind::Start);
1649 }
1650 break;
1652 {
1653
1654
1655
1656
1657
1658 std::string modified_url;
1660 options_op =
new CurlOptionsOp(iter->first, op, target, m_logger, op->GetConnCalloutFunc());
1661 std::shared_ptr<CurlOperation> new_op(options_op);
1662 auto curl = queue.GetHandle();
1663 if (curl == nullptr) {
1664 m_logger->Debug(
kLogXrdClHttp,
"Unable to allocate a curl handle");
1666 keep_handle = false;
1667 options_op = nullptr;
1668 break;
1669 }
1670 OpRecord(*new_op, OpKind::Start);
1671 try {
1672 auto rv = new_op->Setup(curl, *this);
1673 if (!rv) {
1674 m_logger->Debug(
kLogXrdClHttp,
"Unable to configure a curl handle for OPTIONS");
1675 keep_handle = false;
1676 options_op = nullptr;
1677 break;
1678 }
1679 } catch (...) {
1680 m_logger->Debug(
kLogXrdClHttp,
"Unable to setup the curl handle for the OPTIONS operation");
1681 new_op->Fail(
XrdCl::errInternal, ENOMEM,
"Failed to setup the curl handle for the OPTIONS operation");
1682 OpRecord(*new_op, OpKind::Error);
1683 keep_handle = false;
1684 break;
1685 }
1686 new_op->SetContinueQueue(m_continue_queue);
1687 m_op_map[curl] = {new_op, std::chrono::system_clock::now()};
1688 auto mres = curl_multi_add_handle(multi_handle, curl);
1689 if (mres != CURLM_OK) {
1690 m_logger->Debug(
kLogXrdClHttp,
"Unable to add OPTIONS operation to the curl multi-handle: %s", curl_multi_strerror(mres));
1691 op->Fail(
XrdCl::errInternal, mres,
"Unable to add OPTIONS operation to the curl multi-handle");
1692 OpRecord(*new_op, OpKind::Error);
1693 break;
1694 }
1695 running_handles += 1;
1696 m_logger->Debug(
kLogXrdClHttp,
"Invoking the OPTIONS operation before redirect to %s", target.c_str());
1697
1698
1699 keep_handle = true;
1700 }
1701 }
1702 int callout_socket = op->WaitSocket();
1703 if ((waiting_on_callout = callout_socket >= 0)) {
1704 auto expiry = time(nullptr) + 20;
1705 m_logger->Debug(
kLogXrdClHttp,
"Creating a callout wait request on socket %d", callout_socket);
1706 broker_reqs[callout_socket] = {iter->first, expiry};
1707 m_conncall_req.fetch_add(1, std::memory_order_relaxed);
1708 }
1709 } else if (options_op) {
1710
1712 }
1713 if (keep_handle) {
1714 curl_multi_remove_handle(multi_handle, iter->first);
1715 if (!waiting_on_callout && !options_op) {
1716 curl_multi_add_handle(multi_handle, iter->first);
1717 }
1718 } else if (!options_op) {
1719
1720
1721
1722
1723 curl_multi_remove_handle(multi_handle, iter->first);
1724 op->Success();
1725 if (op->IsDone()) {
1726 op->ReleaseHandle();
1727
1728 queue.RecycleHandle(iter->first);
1729 } else {
1730
1731
1732
1733
1734 auto next_res = curl_multi_add_handle(multi_handle, iter->first);
1735 if (next_res == CURLM_OK) {
1736 keep_handle = true;
1737 OpRecord(*op, OpKind::Start);
1738 } else {
1740 "Unable to add the next operation request to the curl multi-handle");
1741 OpRecord(*op, OpKind::Error);
1742 op->ReleaseHandle();
1743 queue.RecycleHandle(iter->first);
1744 }
1745 }
1746 }
1747 }
1748 } else if (res == CURLE_COULDNT_CONNECT && op->UseConnectionCallout() && !op->GetTriedBoker()) {
1749
1750
1751 keep_handle = true;
1752 op->SetTriedBoker();
1753 std::string err;
1754 int wait_socket = -1;
1755 if (!op->StartConnectionCallout(err) || (wait_socket=op->WaitSocket()) == -1) {
1756 m_logger->Error(
kLogXrdClHttp,
"Failed to start broker-based connection: %s", err.c_str());
1757 op->ReleaseHandle();
1758 keep_handle = false;
1759 } else {
1760 curl_multi_remove_handle(multi_handle, iter->first);
1761 auto expiry = time(nullptr) + 20;
1762 m_logger->Debug(
kLogXrdClHttp,
"Curl operation requires a new TCP socket; waiting for callout to respond on socket %d", wait_socket);
1763 broker_reqs[wait_socket] = {iter->first, expiry};
1764 m_conncall_req.fetch_add(1, std::memory_order_relaxed);
1765 }
1766 } else {
1767 if (res == CURLE_ABORTED_BY_CALLBACK || res == CURLE_WRITE_ERROR) {
1768
1769
1770 switch (op->GetError()) {
1772#ifdef HAVE_XPROTOCOL_TIMEREXPIRED
1774#else
1776#endif
1777 OpRecord(*op, OpKind::Error);
1778 break;
1780 auto [ecode,
emsg] = op->GetCallbackError();
1782 OpRecord(*op, OpKind::Error);
1783 break;
1784 }
1787 OpRecord(*op, op->IsPaused() ? OpKind::ClientTimeout : OpKind::ServerTimeout);
1788 break;
1791 OpRecord(*op, OpKind::ServerTimeout);
1792 break;
1795 OpRecord(*op, OpKind::ClientTimeout);
1796 break;
1799 OpRecord(*op, OpKind::ServerTimeout);
1800 break;
1802 op->Fail(
XrdCl::errInternal, 0,
"Operation was aborted without recording an abort reason");
1803 OpRecord(*op, OpKind::Error);
1804 break;
1805 };
1806 CurlOptionsOp *options_op = nullptr;
1807 if ((options_op = dynamic_cast<CurlOptionsOp*>(op.get())) != nullptr) {
1809 bool parent_op_failed = false;
1810 if (parent_op->IsRedirect()) {
1811 std::string target;
1814 if (iter != m_op_map.end()) {
1815 OpRecord(*iter->second.first, OpKind::Error);
1817 m_op_map.erase(iter);
1818 running_handles -= 1;
1819 }
1820 parent_op_failed = true;
1821 } else {
1822 OpRecord(*parent_op, OpKind::Start);
1823 }
1824 } else {
1825 OpRecord(*parent_op, OpKind::Start);
1826 }
1827 if (!parent_op_failed){
1829 }
1830 }
1831 } else {
1833 const auto curl_err = op->GetCurlErrorMessage();
1834 const char *curl_easy_err = curl_easy_strerror(res);
1835 const std::string fail_err = !curl_err.empty() ? curl_err : curl_easy_err;
1836 m_logger->Debug(
kLogXrdClHttp,
"Curl generated an error: %s (%d)", fail_err.c_str(), res);
1837 op->Fail(xrdCode.first, xrdCode.second, fail_err);
1838 OpRecord(*op, OpKind::Error);
1839 CurlOptionsOp *options_op = nullptr;
1840 if ((options_op = dynamic_cast<CurlOptionsOp*>(op.get())) != nullptr) {
1842 bool parent_op_failed = false;
1843 if (parent_op->IsRedirect()) {
1844 std::string target;
1847 if (iter != m_op_map.end()) {
1848 OpRecord(*iter->second.first, OpKind::Error);
1850 m_op_map.erase(iter);
1851 running_handles -= 1;
1852 }
1853 parent_op_failed = true;
1854 }
1855 }
1856 if (!parent_op_failed){
1858 }
1859 }
1860 }
1861 op->ReleaseHandle();
1862 }
1863 if (!keep_handle) {
1864 curl_multi_remove_handle(multi_handle, iter->first);
1865 if (res != CURLE_OK) {
1866 curl_easy_cleanup(iter->first);
1867 }
1868 for (auto &req : broker_reqs) {
1869 if (req.second.curl == iter->first) {
1870 m_logger->Warning(
kLogXrdClHttp,
"Curl handle finished while a broker operation was outstanding");
1871 m_conncall_errors.fetch_add(1, std::memory_order_relaxed);
1872 }
1873 }
1874 m_op_map.erase(iter);
1875 running_handles -= 1;
1876 }
1877 }
1878 } while (msg);
1879 }
1880
1881 for (auto map_entry : m_op_map) {
1882 if (mres) {
1884 OpRecord(*map_entry.second.first, OpKind::Error);
1885 }
1886 if (multi_handle && map_entry.first) curl_multi_remove_handle(multi_handle, map_entry.first);
1887 }
1888
1889 m_queue->ReleaseHandles();
1890 curl_multi_cleanup(multi_handle);
1891}
std::pair< uint16_t, uint32_t > CurlCodeConvert(CURLcode res)
int emsg(int rc, char *msg)
static void CleanupDnsCache()
CurlOptionsOp(CURL *curl, std::shared_ptr< CurlOperation > op, const std::string &url, XrdCl::Log *log, CreateConnCalloutType callout)
std::shared_ptr< CurlOperation > GetOperation() const
CURL * GetParentCurlHandle() const
void ReleaseHandle() override
void Fail(uint16_t errCode, uint32_t errNum, const std::string &) override
HandlerQueue(unsigned max_pending_ops)
static std::string_view GetUrlKey(const std::string &url, std::string &modified_url)
bool GetInt(const std::string &key, int &value)
std::pair< uint16_t, uint32_t > HTTPStatusConvert(unsigned status)
bool HTTPStatusIsError(unsigned status)
const uint64_t kLogXrdClHttp
const uint16_t errErrorResponse
const uint16_t errOperationExpired
const uint16_t errInternal
Internal error.
const uint16_t errConnectionError