XRootD
Loading...
Searching...
No Matches
XrdClHttp::HandlerQueue Class Reference

#include <XrdClHttpUtil.hh>

Collaboration diagram for XrdClHttp::HandlerQueue:

Public Member Functions

 HandlerQueue (unsigned max_pending_ops)
std::shared_ptr< CurlOperation > Consume (std::chrono::steady_clock::duration)
void Expire ()
CURL * GetHandle ()
int PollFD () const
void Produce (std::shared_ptr< CurlOperation > handler)
void RecycleHandle (CURL *)
void ReleaseHandles ()
void Shutdown ()
std::shared_ptr< CurlOperation > TryConsume ()

Static Public Member Functions

static unsigned GetDefaultMaxPendingOps ()
static std::string GetMonitoringJson ()

Detailed Description

HandlerQueue is a deque of curl operations that need to be performed. The object is thread safe and can be waited on via poll().

The fact that it's poll'able is necessary because the multi-curl driver thread is based on polling FD's

Definition at line 191 of file XrdClHttpUtil.hh.

Constructor & Destructor Documentation

◆ HandlerQueue()

HandlerQueue::HandlerQueue ( unsigned max_pending_ops)

Definition at line 684 of file XrdClHttpUtil.cc.

684 :
685 m_max_pending_ops(max_pending_ops)
686{
687 int filedes[2];
688 auto result = pipe(filedes);
689 if (result == -1) {
690 throw std::system_error(errno, std::generic_category(),
691 "failed to create HTTP worker pipe");
692 }
693 if (fcntl(filedes[0], F_SETFL, O_NONBLOCK | O_CLOEXEC) == -1 || fcntl(filedes[1], F_SETFL, O_NONBLOCK | O_CLOEXEC) == -1) {
694 const int error = errno;
695 close(filedes[0]);
696 close(filedes[1]);
697 throw std::system_error(error, std::generic_category(),
698 "failed to configure HTTP worker pipe");
699 }
700 m_read_fd = filedes[0];
701 m_write_fd = filedes[1];
702};
#define close(a)
Definition XrdPosix.hh:48

References close.

Referenced by XrdClHttp::CurlWorker::Run().

Here is the caller graph for this function:

Member Function Documentation

◆ Consume()

std::shared_ptr< CurlOperation > HandlerQueue::Consume ( std::chrono::steady_clock::duration dur)

Definition at line 926 of file XrdClHttpUtil.cc.

927{
928 std::unique_lock<std::mutex> lk(m_mutex);
929 m_consumer_cv.wait_for(lk, dur, [&]{return m_ops.size() > 0 || m_shutdown;});
930 if (m_shutdown || m_ops.empty()) {
931 return {};
932 }
933
934 std::shared_ptr<CurlOperation> result = m_ops.front();
935 m_ops.pop_front();
936
937 char ready[1];
938 while (true) {
939 auto result = read(m_read_fd, ready, 1);
940 if (result == -1) {
941 if (errno == EINTR) {
942 continue;
943 } else if (errno == EAGAIN || errno == EWOULDBLOCK) {
944 // This should never happen, but if it does, just continue
945 // as if we successfully read the byte.
946 break;
947 }
948 throw std::system_error(errno, std::generic_category(),
949 "failed to consume HTTP worker notification");
950 }
951 break;
952 }
953
954 lk.unlock();
955 m_producer_cv.notify_one();
956 m_ops_consumed.fetch_add(1, std::memory_order_relaxed);
957
958 return result;
959}
#define read(a, b, c)
Definition XrdPosix.hh:86

References read.

◆ Expire()

void HandlerQueue::Expire ( )

Definition at line 824 of file XrdClHttpUtil.cc.

825{
826 std::unique_lock<std::mutex> lk(m_mutex);
827 auto now = std::chrono::steady_clock::now();
828
829 // Iterate through the paused transfers, checking if they are done.
830 for (auto &op : m_ops) {
831 if (!op->IsPaused()) continue;
832
833 if (op->TransferStalled(0, now)) {
834 op->ContinueHandle();
835 }
836 }
837
838 std::vector<decltype(m_ops)::value_type> expired_ops;
839 unsigned expired_count = 0;
840 auto it = std::remove_if(m_ops.begin(), m_ops.end(),
841 [&](const std::shared_ptr<CurlOperation> &handler) {
842 auto expired = handler->GetOperationExpiry() < now;
843 if (expired) {
844 expired_ops.push_back(handler);
845 expired_count++;
846 }
847 return expired;
848 });
849 m_ops.erase(it, m_ops.end());
850
851 // The contents of our pipe and the in-memory queue are now off by expired_count.
852 // Read exactly that many bytes from the pipe and throw them away.
853 char throwaway[64];
854 unsigned bytes_to_read = expired_count;
855 while (bytes_to_read > 0) {
856 size_t chunk = std::min<size_t>(sizeof(throwaway), bytes_to_read);
857 ssize_t n = read(m_read_fd, throwaway, chunk);
858 if (n > 0) {
859 bytes_to_read -= n;
860 } else if (n == -1) {
861 if (errno == EINTR) {
862 continue;
863 } else {
864 // EWOULDBLOCK is a possibility if there's a synchronization error;
865 // for now, just continue on as if we were successful in reading out
866 // the missing bytes
867 break;
868 }
869 } else {
870 break;
871 }
872 }
873
874 // Note: the failure handler may trigger new operations submitted to the queue
875 // (which requires the lock to be held) such as a prefetch operation that gets split
876 // into multiple sub-operations.
877 //
878 // Thus, we must unlock the mutex protecting the queue and avoid touching the shared state of
879 // m_ops.
880 lk.unlock();
881 for (auto &handler : expired_ops) {
882 if (handler) handler->Fail(XrdCl::errOperationExpired, 0, "Operation expired while in queue");
883 }
884}
const uint16_t errOperationExpired

◆ GetDefaultMaxPendingOps()

unsigned XrdClHttp::HandlerQueue::GetDefaultMaxPendingOps ( )
inlinestatic

Definition at line 218 of file XrdClHttpUtil.hh.

218{return m_default_max_pending_ops;}

◆ GetHandle()

CURL * HandlerQueue::GetHandle ( )

Definition at line 808 of file XrdClHttpUtil.cc.

808 {
809 if (m_handles.size()) {
810 auto result = m_handles.back();
811 m_handles.pop_back();
812 return result;
813 }
814
815 return ::GetHandle(EnableCurlHeaderDump());
816}

◆ GetMonitoringJson()

std::string HandlerQueue::GetMonitoringJson ( )
static

Definition at line 962 of file XrdClHttpUtil.cc.

963{
964 auto consumed = m_ops_consumed.load(std::memory_order_relaxed);
965 auto produced = m_ops_produced.load(std::memory_order_relaxed);
966 return "{"
967 "\"produced\":" + std::to_string(produced) + ","
968 "\"consumed\":" + std::to_string(consumed) + ","
969 "\"pending\":" + std::to_string(produced - consumed) + ","
970 "\"rejected\":" + std::to_string(m_ops_rejected.load(std::memory_order_relaxed)) +
971 "}";
972}

◆ PollFD()

int XrdClHttp::HandlerQueue::PollFD ( ) const
inline

Definition at line 200 of file XrdClHttpUtil.hh.

200{return m_read_fd;}

◆ Produce()

void HandlerQueue::Produce ( std::shared_ptr< CurlOperation > handler)

Definition at line 887 of file XrdClHttpUtil.cc.

888{
889 auto handler_expiry = handler->GetOperationExpiry();
890 std::unique_lock<std::mutex> lk{m_mutex};
891 m_producer_cv.wait_until(lk,
892 handler_expiry,
893 [&]{return m_ops.size() < m_max_pending_ops;}
894 );
895 if (std::chrono::steady_clock::now() > handler_expiry) {
896 lk.unlock();
897 handler->Fail(XrdCl::errOperationExpired, 0, "Operation expired while waiting for worker");
898 m_ops_rejected.fetch_add(1, std::memory_order_relaxed);
899 return;
900 }
901
902 m_ops.push_back(handler);
903 char ready[] = "1";
904 while (true) {
905 auto result = write(m_write_fd, ready, 1);
906 if (result == -1) {
907 if (errno == EINTR) {
908 continue;
909 } else if (errno == EAGAIN || errno == EWOULDBLOCK) {
910 // This should never happen, but if it does, just continue
911 // as if we successfully wrote the notification to the pipe.
912 break;
913 }
914 throw std::system_error(errno, std::generic_category(),
915 "failed to notify HTTP worker");
916 }
917 break;
918 }
919
920 lk.unlock();
921 m_consumer_cv.notify_one();
922 m_ops_produced.fetch_add(1, std::memory_order_relaxed);
923}
#define write(a, b, c)
Definition XrdPosix.hh:121

References XrdCl::errOperationExpired, and write.

Referenced by XrdClHttp::CurlReadOp::Continue(), and XrdClHttp::stat_size().

Here is the caller graph for this function:

◆ RecycleHandle()

void HandlerQueue::RecycleHandle ( CURL * curl)

Definition at line 819 of file XrdClHttpUtil.cc.

819 {
820 m_handles.push_back(curl);
821}

◆ ReleaseHandles()

void HandlerQueue::ReleaseHandles ( )

Definition at line 1019 of file XrdClHttpUtil.cc.

1020{
1021 for (auto handle : m_handles) {
1022 curl_easy_cleanup(handle);
1023 }
1024 m_handles.clear();
1025}

◆ Shutdown()

void HandlerQueue::Shutdown ( )

Definition at line 1011 of file XrdClHttpUtil.cc.

1012{
1013 std::unique_lock lock(m_mutex);
1014 m_shutdown = true;
1015 m_consumer_cv.notify_all();
1016}

◆ TryConsume()

std::shared_ptr< CurlOperation > HandlerQueue::TryConsume ( )

Definition at line 975 of file XrdClHttpUtil.cc.

976{
977 std::unique_lock<std::mutex> lk(m_mutex);
978 if (m_ops.size() == 0) {
979 std::shared_ptr<CurlOperation> result;
980 return result;
981 }
982
983 std::shared_ptr<CurlOperation> result = m_ops.front();
984 m_ops.pop_front();
985
986 char ready[1];
987 while (true) {
988 auto result = read(m_read_fd, ready, 1);
989 if (result == -1) {
990 if (errno == EINTR) {
991 continue;
992 } else if (errno == EAGAIN || errno == EWOULDBLOCK) {
993 // This should never happen, but if it does, just continue
994 // as if we successfully read the byte.
995 break;
996 }
997 throw std::system_error(errno, std::generic_category(),
998 "failed to consume HTTP worker notification");
999 }
1000 break;
1001 }
1002
1003 lk.unlock();
1004 m_producer_cv.notify_one();
1005 m_ops_consumed.fetch_add(1, std::memory_order_relaxed);
1006
1007 return result;
1008}

References read.


The documentation for this class was generated from the following files: