diff --git a/configure.ac b/configure.ac index ce162123..34a92101 100644 --- a/configure.ac +++ b/configure.ac @@ -26,9 +26,17 @@ TORRENT_WITHOUT_NCURSESW() TORRENT_WITHOUT_STATVFS() TORRENT_WITHOUT_STATFS() -PKG_CHECK_MODULES(STUFF, sigc++-2.0 libcurl >= 7.12.0 libtorrent >= 0.12.2, - CXXFLAGS="$CXXFLAGS $STUFF_CFLAGS"; - LIBS="$LIBS $STUFF_LIBS") +PKG_CHECK_MODULES(sigc, sigc++-2.0, + CXXFLAGS="$CXXFLAGS $sigc_CFLAGS"; + LIBS="$LIBS $sigc_LIBS") + +PKG_CHECK_MODULES(libcurl, libcurl >= 7.15.4, + CXXFLAGS="$CXXFLAGS $libcurl_CFLAGS"; + LIBS="$LIBS $libcurl_LIBS") + +PKG_CHECK_MODULES(libtorrent, libtorrent >= 0.12.2, + CXXFLAGS="$CXXFLAGS $libtorrent_CFLAGS"; + LIBS="$LIBS $libtorrent_LIBS") AC_LANG_PUSH(C++) TORRENT_WITH_XMLRPC_C diff --git a/src/command_network.cc b/src/command_network.cc index 551fc9cf..09e4915d 100644 --- a/src/command_network.cc +++ b/src/command_network.cc @@ -291,7 +291,7 @@ apply_xmlrpc_dialect(const std::string& arg) { void initialize_command_network() { torrent::ConnectionManager* cm = torrent::connection_manager(); - core::CurlStack* httpStack = control->core()->get_poll_manager()->get_http_stack(); + core::CurlStack* httpStack = control->core()->http_stack(); ADD_VARIABLE_BOOL("use_udp_trackers", true); diff --git a/src/control.cc b/src/control.cc index d20be124..649fb11b 100644 --- a/src/control.cc +++ b/src/control.cc @@ -107,7 +107,7 @@ Control::initialize() { display::Window::slot_unschedule(rak::make_mem_fun(m_display, &display::Manager::unschedule)); display::Window::slot_adjust(rak::make_mem_fun(m_display, &display::Manager::adjust_layout)); - m_core->get_poll_manager()->get_http_stack()->set_user_agent(USER_AGENT); + m_core->http_stack()->set_user_agent(USER_AGENT); m_core->initialize_second(); m_core->listen_open(); diff --git a/src/core/Makefile.am b/src/core/Makefile.am index 430617a0..e1bf1645 100644 --- a/src/core/Makefile.am +++ b/src/core/Makefile.am @@ -3,6 +3,8 @@ noinst_LIBRARIES = libsub_core.a libsub_core_a_SOURCES = \ curl_get.cc \ curl_get.h \ + curl_socket.cc \ + curl_socket.h \ curl_stack.cc \ curl_stack.h \ dht_manager.cc \ diff --git a/src/core/curl_stack.cc b/src/core/curl_stack.cc index 4522cf09..021a3383 100644 --- a/src/core/curl_stack.cc +++ b/src/core/curl_stack.cc @@ -43,6 +43,7 @@ #include "rak/functional.h" #include "curl_get.h" +#include "curl_socket.h" #include "curl_stack.h" namespace core { @@ -51,6 +52,15 @@ CurlStack::CurlStack() : m_handle((void*)curl_multi_init()), m_active(0), m_maxActive(32) { + + m_taskTimeout.set_slot(rak::mem_fn(this, &CurlStack::receive_timeout)); + + curl_multi_setopt((CURLM*)m_handle, CURLMOPT_TIMERDATA, this); + curl_multi_setopt((CURLM*)m_handle, CURLMOPT_SOCKETDATA, this); +#if (LIBCURL_VERSION_NUM >= 0x071000) + curl_multi_setopt((CURLM*)m_handle, CURLMOPT_TIMERFUNCTION, &CurlStack::set_timeout); +#endif + curl_multi_setopt((CURLM*)m_handle, CURLMOPT_SOCKETFUNCTION, &CurlSocket::receive_socket); } CurlStack::~CurlStack() { @@ -58,6 +68,7 @@ CurlStack::~CurlStack() { front()->close(); curl_multi_cleanup((CURLM*)m_handle); + priority_queue_erase(&taskScheduler, &m_taskTimeout); } CurlGet* @@ -65,16 +76,31 @@ CurlStack::new_object() { return new CurlGet(this); } +CurlSocket* +CurlStack::new_socket(int fd) { + CurlSocket* socket = new CurlSocket(fd, this); + curl_multi_assign((CURLM*)m_handle, fd, socket); + return socket; +} + void -CurlStack::perform() { +CurlStack::receive_action(CurlSocket* socket, int events) { CURLMcode code; do { int count; - code = curl_multi_perform((CURLM*)m_handle, &count); +#if (LIBCURL_VERSION_NUM >= 0x071003) + code = curl_multi_socket_action((CURLM*)m_handle, socket != NULL ? socket->file_descriptor() : CURL_SOCKET_TIMEOUT, events, &count); +#else + code = curl_multi_socket((CURLM*)m_handle, socket != NULL ? socket->file_descriptor() : CURL_SOCKET_TIMEOUT, &count); +#endif if (code > 0) - throw torrent::internal_error("Error calling curl_multi_perform."); + throw torrent::internal_error("Error calling curl_multi_socket_action."); + + // Socket might be removed when cleaning handles below, future calls should not use it. + socket = NULL; + events = 0; if ((unsigned int)count != size()) { // Done with some handles. @@ -83,31 +109,29 @@ CurlStack::perform() { while ((msg = curl_multi_info_read((CURLM*)m_handle, &t)) != NULL) { if (msg->msg != CURLMSG_DONE) - throw torrent::internal_error("CurlStack::perform() msg->msg != CURLMSG_DONE."); + throw torrent::internal_error("CurlStack::receive_action() msg->msg != CURLMSG_DONE."); iterator itr = std::find_if(begin(), end(), rak::equal(msg->easy_handle, std::mem_fun(&CurlGet::handle))); if (itr == end()) throw torrent::internal_error("Could not find CurlGet with the right easy_handle."); - + if (msg->data.result == CURLE_OK) (*itr)->signal_done().emit(); else (*itr)->signal_failed().emit(curl_easy_strerror(msg->data.result)); } + + if (empty()) + priority_queue_erase(&taskScheduler, &m_taskTimeout); } } while (code == CURLM_CALL_MULTI_PERFORM); } -unsigned int -CurlStack::fdset(fd_set* readfds, fd_set* writefds, fd_set* exceptfds) { - int maxFd = 0; - - if (curl_multi_fdset((CURLM*)m_handle, readfds, writefds, exceptfds, &maxFd) != 0) - throw torrent::internal_error("Error calling curl_multi_fdset."); - - return std::max(maxFd, 0); +void +CurlStack::receive_timeout() { + receive_action(NULL, 0); } void @@ -179,4 +203,14 @@ CurlStack::global_cleanup() { curl_global_cleanup(); } +int +CurlStack::set_timeout(void* handle, long timeout_ms, void* userp) { + CurlStack* stack = (CurlStack*)userp; + + priority_queue_erase(&taskScheduler, &stack->m_taskTimeout); + priority_queue_insert(&taskScheduler, &stack->m_taskTimeout, cachedTime + rak::timer::from_milliseconds(timeout_ms)); + + return 0; +} + } diff --git a/src/core/curl_stack.h b/src/core/curl_stack.h index b1adf4fc..e199379c 100644 --- a/src/core/curl_stack.h +++ b/src/core/curl_stack.h @@ -41,9 +41,12 @@ #include #include +#include "rak/priority_queue_default.h" + namespace core { class CurlGet; +class CurlSocket; // By using a deque instead of vector we allow for cheaper removal of // the oldest elements, those that will be first in the in the @@ -80,11 +83,7 @@ class CurlStack : std::deque { ~CurlStack(); CurlGet* new_object(); - - void perform(); - - // TODO: Set fd_set's only once? - unsigned int fdset(fd_set* readfds, fd_set* writefds, fd_set* exceptfds); + CurlSocket* new_socket(int fd); unsigned int active() const { return m_active; } unsigned int max_active() const { return m_maxActive; } @@ -108,6 +107,10 @@ class CurlStack : std::deque { static void global_init(); static void global_cleanup(); + void receive_action(CurlSocket* socket, int type); + + static int set_timeout(void* handle, long timeout_ms, void* userp); + protected: void add_get(CurlGet* get); void remove_get(CurlGet* get); @@ -116,11 +119,15 @@ class CurlStack : std::deque { CurlStack(const CurlStack&); void operator = (const CurlStack&); + void receive_timeout(); + void* m_handle; unsigned int m_active; unsigned int m_maxActive; + rak::priority_item m_taskTimeout; + std::string m_userAgent; std::string m_httpProxy; std::string m_bindAddress; diff --git a/src/core/manager.cc b/src/core/manager.cc index 5ec5feb4..2fba8ad8 100644 --- a/src/core/manager.cc +++ b/src/core/manager.cc @@ -165,6 +165,7 @@ Manager::Manager() : m_downloadList = new DownloadList(); m_fileStatusCache = new FileStatusCache(); m_httpQueue = new HttpQueue(); + m_httpStack = new CurlStack(); } Manager::~Manager() { @@ -205,8 +206,8 @@ Manager::initialize_first() { // Most of this should be possible to move out. void Manager::initialize_second() { - torrent::Http::set_factory(sigc::mem_fun(m_pollManager->get_http_stack(), &CurlStack::new_object)); - m_httpQueue->slot_factory(sigc::mem_fun(m_pollManager->get_http_stack(), &CurlStack::new_object)); + torrent::Http::set_factory(sigc::mem_fun(m_httpStack, &CurlStack::new_object)); + m_httpQueue->slot_factory(sigc::mem_fun(m_httpStack, &CurlStack::new_object)); CurlStack::global_init(); @@ -229,6 +230,8 @@ Manager::cleanup() { // here before the torrent::* objects are deleted. torrent::cleanup(); + + delete m_httpStack; CurlStack::global_cleanup(); delete m_pollManager; @@ -304,7 +307,7 @@ Manager::set_bind_address(const std::string& addr) { torrent::connection_manager()->set_bind_address(ai->address()->c_sockaddr()); } - m_pollManager->get_http_stack()->set_bind_address(!ai->address()->is_address_any() ? ai->address()->address_str() : std::string()); + m_httpStack->set_bind_address(!ai->address()->is_address_any() ? ai->address()->address_str() : std::string()); rak::address_info::free_address_info(ai); diff --git a/src/core/manager.h b/src/core/manager.h index ad823923..43c10156 100644 --- a/src/core/manager.h +++ b/src/core/manager.h @@ -75,6 +75,7 @@ public: FileStatusCache* file_status_cache() { return m_fileStatusCache; } HttpQueue* http_queue() { return m_httpQueue; } + CurlStack* http_stack() { return m_httpStack; } View* hashing_view() { return m_hashingView; } void set_hashing_view(View* v); @@ -131,6 +132,7 @@ private: DownloadStore* m_downloadStore; FileStatusCache* m_fileStatusCache; HttpQueue* m_httpQueue; + CurlStack* m_httpStack; View* m_hashingView; diff --git a/src/core/poll_manager.cc b/src/core/poll_manager.cc index d5bf579d..7e9600f3 100644 --- a/src/core/poll_manager.cc +++ b/src/core/poll_manager.cc @@ -48,47 +48,10 @@ PollManager::PollManager(torrent::Poll* poll) : if (m_poll == NULL) throw std::logic_error("PollManager::PollManager(...) received poll == NULL"); - -#if defined USE_VARIABLE_FDSET - m_setSize = m_poll->open_max() / 8; - m_readSet = (fd_set*)new char[m_setSize]; - m_writeSet = (fd_set*)new char[m_setSize]; - m_errorSet = (fd_set*)new char[m_setSize]; - - std::memset(m_readSet, 0, m_setSize); - std::memset(m_writeSet, 0, m_setSize); - std::memset(m_errorSet, 0, m_setSize); -#else - if (m_poll->open_max() > FD_SETSIZE) - throw std::logic_error("PollManager::PollManager(...) received a max open sockets >= FD_SETSIZE, but USE_VARIABLE_FDSET was not defined"); - - m_setSize = FD_SETSIZE / 8; - m_readSet = new fd_set; - m_writeSet = new fd_set; - m_errorSet = new fd_set; - - FD_ZERO(m_readSet); - FD_ZERO(m_writeSet); - FD_ZERO(m_errorSet); -#endif - - // Call this so curl has valid fd_set pointers if curl_multi_perform - // is called before it gets set when polling. - m_httpStack.fdset(m_readSet, m_writeSet, m_errorSet); } PollManager::~PollManager() { delete m_poll; - -#if defined USE_VARIABLE_FDSET - delete [] m_readSet; - delete [] m_writeSet; - delete [] m_errorSet; -#else - delete m_readSet; - delete m_writeSet; - delete m_errorSet; -#endif } void diff --git a/src/core/poll_manager.h b/src/core/poll_manager.h index 61dfd380..1da2f432 100644 --- a/src/core/poll_manager.h +++ b/src/core/poll_manager.h @@ -37,7 +37,6 @@ #ifndef RTORRENT_CORE_POLL_MANAGER_H #define RTORRENT_CORE_POLL_MANAGER_H -#include #include #include #include @@ -58,7 +57,6 @@ public: unsigned int get_open_max() const { return m_poll->open_max(); } - CurlStack* get_http_stack() { return &m_httpStack; } torrent::Poll* get_torrent_poll() { return m_poll; } virtual void poll(rak::timer timeout) = 0; @@ -70,12 +68,6 @@ protected: void check_error(); torrent::Poll* m_poll; - CurlStack m_httpStack; - - unsigned int m_setSize; - fd_set* m_readSet; - fd_set* m_writeSet; - fd_set* m_errorSet; }; } diff --git a/src/core/poll_manager_epoll.cc b/src/core/poll_manager_epoll.cc index c6d7db78..c09c85f2 100644 --- a/src/core/poll_manager_epoll.cc +++ b/src/core/poll_manager_epoll.cc @@ -67,44 +67,6 @@ PollManagerEPoll::poll(rak::timer timeout) { torrent::perform(); timeout = std::min(timeout, rak::timer(torrent::next_timeout())) + 1000; - if (!m_httpStack.empty()) { - // When we're using libcurl we need to use select, but as this is - // inefficient we try avoiding it whenever possible. -#if defined USE_VARIABLE_FDSET - std::memset(m_readSet, 0, m_setSize); - std::memset(m_writeSet, 0, m_setSize); - std::memset(m_errorSet, 0, m_setSize); -#else - FD_ZERO(m_readSet); - FD_ZERO(m_writeSet); - FD_ZERO(m_errorSet); -#endif - FD_SET(static_cast(m_poll)->file_descriptor(), m_readSet); - - unsigned int maxFd = std::max((unsigned int)static_cast(m_poll)->file_descriptor(), - m_httpStack.fdset(m_readSet, m_writeSet, m_errorSet)); - - timeval t = timeout.tval(); - - if (select(maxFd + 1, m_readSet, m_writeSet, m_errorSet, &t) == -1) - return check_error(); - - m_httpStack.perform(); - - if (!FD_ISSET(static_cast(m_poll)->file_descriptor(), m_readSet)) { - // Need to call perform here so that scheduled task get done - // even if there's no socket events outside of the http stuff. - torrent::perform(); - return; - } - - // Clear the timeout since we've already used it in the select call. - timeout = rak::timer(); - } - - // Yes, below is how much code really *should* have been in this - // function. ;) - if (static_cast(m_poll)->poll((timeout.usec() + 999) / 1000) == -1) return check_error(); diff --git a/src/core/poll_manager_kqueue.cc b/src/core/poll_manager_kqueue.cc index 75313434..641b4d51 100644 --- a/src/core/poll_manager_kqueue.cc +++ b/src/core/poll_manager_kqueue.cc @@ -68,44 +68,6 @@ PollManagerKQueue::poll(rak::timer timeout) { torrent::perform(); timeout = std::min(timeout, rak::timer(torrent::next_timeout())) + 1000; - if (!m_httpStack.empty()) { - // When we're using libcurl we need to use select, but as this is - // inefficient we try avoiding it whenever possible. -#if defined USE_VARIABLE_FDSET - std::memset(m_readSet, 0, m_setSize); - std::memset(m_writeSet, 0, m_setSize); - std::memset(m_errorSet, 0, m_setSize); -#else - FD_ZERO(m_readSet); - FD_ZERO(m_writeSet); - FD_ZERO(m_errorSet); -#endif - FD_SET(static_cast(m_poll)->file_descriptor(), m_readSet); - - unsigned int maxFd = std::max((unsigned int)static_cast(m_poll)->file_descriptor(), - m_httpStack.fdset(m_readSet, m_writeSet, m_errorSet)); - - timeval t = timeout.tval(); - - if (select(maxFd + 1, m_readSet, m_writeSet, m_errorSet, &t) == -1) - return check_error(); - - m_httpStack.perform(); - - if (!FD_ISSET(static_cast(m_poll)->file_descriptor(), m_readSet)) { - // Need to call perform here so that scheduled task get done - // even if there's no socket events outside of the http stuff. - torrent::perform(); - return; - } - - // Clear the timeout since we've already used it in the select call. - timeout = rak::timer(); - } - - // Yes, below is how much code really *should* have been in this - // function. ;) - if (static_cast(m_poll)->poll((timeout.usec() + 999) / 1000) == -1) return check_error(); diff --git a/src/core/poll_manager_select.cc b/src/core/poll_manager_select.cc index 853a1a1a..aca87724 100644 --- a/src/core/poll_manager_select.cc +++ b/src/core/poll_manager_select.cc @@ -47,6 +47,31 @@ namespace core { +PollManagerSelect::PollManagerSelect(torrent::Poll* p) : PollManager(p) { +#if defined USE_VARIABLE_FDSET + m_setSize = m_poll->open_max() / 8; + m_readSet = (fd_set*)new char[m_setSize]; + m_writeSet = (fd_set*)new char[m_setSize]; + m_errorSet = (fd_set*)new char[m_setSize]; + + std::memset(m_readSet, 0, m_setSize); + std::memset(m_writeSet, 0, m_setSize); + std::memset(m_errorSet, 0, m_setSize); +#else + if (m_poll->open_max() > FD_SETSIZE) + throw std::logic_error("PollManagerSelect::create(...) received a max open sockets >= FD_SETSIZE, but USE_VARIABLE_FDSET was not defined"); + + m_setSize = FD_SETSIZE / 8; + m_readSet = new fd_set; + m_writeSet = new fd_set; + m_errorSet = new fd_set; + + FD_ZERO(m_readSet); + FD_ZERO(m_writeSet); + FD_ZERO(m_errorSet); +#endif +} + PollManagerSelect* PollManagerSelect::create(int maxOpenSockets) { torrent::PollSelect* p = torrent::PollSelect::create(maxOpenSockets); @@ -58,6 +83,15 @@ PollManagerSelect::create(int maxOpenSockets) { } PollManagerSelect::~PollManagerSelect() { +#if defined USE_VARIABLE_FDSET + delete [] m_readSet; + delete [] m_writeSet; + delete [] m_errorSet; +#else + delete m_readSet; + delete m_writeSet; + delete m_errorSet; +#endif } void @@ -77,17 +111,11 @@ PollManagerSelect::poll(rak::timer timeout) { unsigned int maxFd = static_cast(m_poll)->fdset(m_readSet, m_writeSet, m_errorSet); - if (!m_httpStack.empty()) - maxFd = std::max(maxFd, m_httpStack.fdset(m_readSet, m_writeSet, m_errorSet)); - timeval t = timeout.tval(); if (select(maxFd + 1, m_readSet, m_writeSet, m_errorSet, &t) == -1) return check_error(); - if (!m_httpStack.empty()) - m_httpStack.perform(); - torrent::perform(); static_cast(m_poll)->perform(m_readSet, m_writeSet, m_errorSet); } diff --git a/src/core/poll_manager_select.h b/src/core/poll_manager_select.h index da0540e7..3b9124f7 100644 --- a/src/core/poll_manager_select.h +++ b/src/core/poll_manager_select.h @@ -53,7 +53,12 @@ public: void poll(rak::timer timeout); private: - PollManagerSelect(torrent::Poll* p) : PollManager(p) {} + PollManagerSelect(torrent::Poll* p); + + unsigned int m_setSize; + fd_set* m_readSet; + fd_set* m_writeSet; + fd_set* m_errorSet; }; } diff --git a/src/display/utils.cc b/src/display/utils.cc index 0cd17890..b39f9c26 100644 --- a/src/display/utils.cc +++ b/src/display/utils.cc @@ -296,8 +296,8 @@ print_status_extra(char* first, char* last) { torrent::max_download_unchoked()); first = print_buffer(first, last, " [H %u/%u]", - control->core()->get_poll_manager()->get_http_stack()->active(), - control->core()->get_poll_manager()->get_http_stack()->max_active()); + control->core()->http_stack()->active(), + control->core()->http_stack()->max_active()); first = print_buffer(first, last, " [S %i/%i/%i]", torrent::total_handshakes(),