* Use curl's new polling interface so we don't have to switch to

select() while http transfers are active.


git-svn-id: svn://rakshasa.no/libtorrent/trunk/rtorrent@1065 e378c898-3ddf-0310-93e7-cc216c733640
This commit is contained in:
rakshasa
2008-08-27 08:41:08 +00:00
parent c627d6dc0c
commit be8b59b00b
15 changed files with 124 additions and 156 deletions
+11 -3
View File
@@ -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
+1 -1
View File
@@ -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);
+1 -1
View File
@@ -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();
+2
View File
@@ -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 \
+47 -13
View File
@@ -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;
}
}
+12 -5
View File
@@ -41,9 +41,12 @@
#include <string>
#include <sigc++/functors/slot.h>
#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<CurlGet*> {
~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<CurlGet*> {
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<CurlGet*> {
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;
+6 -3
View File
@@ -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);
+2
View File
@@ -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;
-37
View File
@@ -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
-8
View File
@@ -37,7 +37,6 @@
#ifndef RTORRENT_CORE_POLL_MANAGER_H
#define RTORRENT_CORE_POLL_MANAGER_H
#include <sys/select.h>
#include <rak/timer.h>
#include <sigc++/signal.h>
#include <torrent/poll.h>
@@ -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;
};
}
-38
View File
@@ -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<torrent::PollEPoll*>(m_poll)->file_descriptor(), m_readSet);
unsigned int maxFd = std::max((unsigned int)static_cast<torrent::PollEPoll*>(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<torrent::PollEPoll*>(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<torrent::PollEPoll*>(m_poll)->poll((timeout.usec() + 999) / 1000) == -1)
return check_error();
-38
View File
@@ -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<torrent::PollKQueue*>(m_poll)->file_descriptor(), m_readSet);
unsigned int maxFd = std::max((unsigned int)static_cast<torrent::PollKQueue*>(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<torrent::PollKQueue*>(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<torrent::PollKQueue*>(m_poll)->poll((timeout.usec() + 999) / 1000) == -1)
return check_error();
+34 -6
View File
@@ -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<torrent::PollSelect*>(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<torrent::PollSelect*>(m_poll)->perform(m_readSet, m_writeSet, m_errorSet);
}
+6 -1
View File
@@ -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;
};
}
+2 -2
View File
@@ -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(),