* Limit the number of simultanious http connections.

* Added max_open_http to the man page.


git-svn-id: svn://rakshasa.no/libtorrent/trunk/rtorrent@840 e378c898-3ddf-0310-93e7-cc216c733640
This commit is contained in:
rakshasa
2007-01-02 17:06:08 +00:00
parent de744ad11d
commit 3b38c4ab9c
14 changed files with 141 additions and 134 deletions
+3 -44
View File
@@ -37,9 +37,9 @@
#include "config.h"
#include <iostream>
#include <stdexcept>
#include <curl/curl.h>
#include <curl/easy.h>
#include <torrent/exceptions.h>
#include "curl_get.h"
#include "curl_stack.h"
@@ -54,30 +54,17 @@ curl_get_receive_write(void* data, size_t size, size_t nmemb, void* handle) {
return 0;
}
CurlGet::CurlGet(CurlStack* s) :
m_handle(NULL),
m_stack(s) {
if (m_stack == NULL)
throw std::logic_error("Tried to create CurlGet without a valid CurlStack");
}
CurlGet::~CurlGet() {
close();
}
CurlGet*
CurlGet::new_object(CurlStack* s) {
return new CurlGet(s);
}
void
CurlGet::start() {
if (is_busy())
throw std::logic_error("Tried to call CurlGet::start on a busy object");
throw torrent::internal_error("Tried to call CurlGet::start on a busy object.");
if (m_stream == NULL)
throw std::logic_error("Tried to call CurlGet::start without a valid output stream");
throw torrent::internal_error("Tried to call CurlGet::start without a valid output stream.");
m_handle = curl_easy_init();
@@ -127,32 +114,4 @@ CurlGet::size_total() {
return d;
}
void
CurlGet::set_user_agent(const char* s) {
curl_easy_setopt(m_handle, CURLOPT_USERAGENT, s);
}
void
CurlGet::set_http_proxy(const char* s) {
curl_easy_setopt(m_handle, CURLOPT_PROXY, s);
}
void
CurlGet::set_bind_address(const char* s) {
curl_easy_setopt(m_handle, CURLOPT_INTERFACE, s);
}
void
CurlGet::perform(CURLMsg* msg) {
if (msg->msg != CURLMSG_DONE)
throw std::logic_error("CurlGet::process got CURLMSG that isn't done");
if (msg->data.result == CURLE_OK)
m_signalDone.emit();
else
m_signalFailed.emit(curl_easy_strerror(msg->data.result));
// Do nothing below after emitting the signals.
}
}
+11 -20
View File
@@ -40,47 +40,38 @@
#include <iosfwd>
#include <string>
#include <curl/curl.h>
#include <torrent/http.h>
#include <sigc++/signal.h>
struct CURLMsg;
#include <torrent/http.h>
namespace core {
class CurlStack;
class CurlGet : public torrent::Http {
public:
friend class CurlStack;
CurlGet(CurlStack* s);
public:
CurlGet(CurlStack* s) : m_active(false), m_handle(NULL), m_stack(s) {}
virtual ~CurlGet();
static CurlGet* new_object(CurlStack* s);
void start();
void close();
bool is_busy() { return m_handle; }
bool is_busy() const { return m_handle; }
bool is_active() const { return m_active; }
void set_active(bool a) { m_active = a; }
double size_done();
double size_total();
void set_user_agent(const char* s);
void set_http_proxy(const char* s);
void set_bind_address(const char* s);
CURL* handle() { return m_handle; }
protected:
CURL* handle() { return m_handle; }
void perform(CURLMsg* msg);
private:
private:
CurlGet(const CurlGet&);
void operator = (const CurlGet&);
CURL* m_handle;
bool m_active;
CURL* m_handle;
CurlStack* m_stack;
};
+57 -35
View File
@@ -37,7 +37,6 @@
#include "config.h"
#include <algorithm>
#include <stdexcept>
#include <curl/multi.h>
#include <sigc++/bind.h>
#include <torrent/exceptions.h>
@@ -50,39 +49,51 @@ namespace core {
CurlStack::CurlStack() :
m_handle((void*)curl_multi_init()),
m_size(0) {
m_active(0),
m_maxActive(32) {
}
CurlStack::~CurlStack() {
while (!m_getList.empty())
m_getList.front()->close();
while (!empty())
front()->close();
curl_multi_cleanup((CURLM*)m_handle);
}
CurlGet*
CurlStack::new_object() {
return new CurlGet(this);
}
void
CurlStack::perform() {
int s;
CURLMcode code;
do {
code = curl_multi_perform((CURLM*)m_handle, &s);
int count;
code = curl_multi_perform((CURLM*)m_handle, &count);
if (code > 0)
throw std::runtime_error("Error calling curl_multi_perform");
throw torrent::internal_error("Error calling curl_multi_perform.");
if (s != m_size) {
if ((unsigned int)count != size()) {
// Done with some handles.
int t;
CURLMsg* msg;
while ((msg = curl_multi_info_read((CURLM*)m_handle, &t)) != NULL) {
CurlGetList::iterator itr = std::find_if(m_getList.begin(), m_getList.end(), rak::equal(msg->easy_handle, std::mem_fun(&CurlGet::handle)));
if (msg->msg != CURLMSG_DONE)
throw torrent::internal_error("CurlStack::perform() msg->msg != CURLMSG_DONE.");
if (itr == m_getList.end())
throw std::logic_error("Could not find CurlGet with the right easy_handle");
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.");
(*itr)->perform(msg);
if (msg->data.result == CURLE_OK)
(*itr)->signal_done().emit();
else
(*itr)->signal_failed().emit(curl_easy_strerror(msg->data.result));
}
}
@@ -94,7 +105,7 @@ 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 std::runtime_error("Error calling curl_multi_fdset");
throw torrent::internal_error("Error calling curl_multi_fdset.");
return std::max(maxFd, 0);
}
@@ -102,38 +113,54 @@ CurlStack::fdset(fd_set* readfds, fd_set* writefds, fd_set* exceptfds) {
void
CurlStack::add_get(CurlGet* get) {
if (!m_userAgent.empty())
get->set_user_agent(m_userAgent.c_str());
curl_easy_setopt(get->handle(), CURLOPT_USERAGENT, m_userAgent.c_str());
if (!m_httpProxy.empty())
get->set_http_proxy(m_httpProxy.c_str());
curl_easy_setopt(get->handle(), CURLOPT_PROXY, m_httpProxy.c_str());
if (!m_bindAddress.empty())
get->set_bind_address(m_bindAddress.c_str());
curl_easy_setopt(get->handle(), CURLOPT_INTERFACE, m_bindAddress.c_str());
CURLMcode code;
base_type::push_back(get);
if ((code = curl_multi_add_handle((CURLM*)m_handle, get->handle())) > 0)
throw std::logic_error("curl_multi_add_handle \"" + std::string(curl_multi_strerror(code)));
if (m_active >= m_maxActive)
return;
m_size++;
m_getList.push_back(get);
// Curl ML suggest we need to do perform after adding a handle.
//perform();
m_active++;
get->set_active(true);
if (curl_multi_add_handle((CURLM*)m_handle, get->handle()) > 0)
throw torrent::internal_error("Error calling curl_multi_add_handle.");
}
void
CurlStack::remove_get(CurlGet* get) {
iterator itr = std::find(begin(), end(), get);
if (itr == end())
throw torrent::internal_error("Could not find CurlGet when calling CurlStack::remove.");
base_type::erase(itr);
// The CurlGet object was never activated, so we just skip this one.
if (!get->is_active())
return;
get->set_active(false);
if (curl_multi_remove_handle((CURLM*)m_handle, get->handle()) > 0)
throw std::logic_error("Error calling curl_multi_remove_handle");
throw torrent::internal_error("Error calling curl_multi_remove_handle.");
CurlGetList::iterator itr = std::find(m_getList.begin(), m_getList.end(), get);
if (m_active == m_maxActive &&
(itr = std::find_if(begin(), end(), std::not1(std::mem_fun(&CurlGet::is_active)))) != end()) {
(*itr)->set_active(true);
if (itr == m_getList.end())
throw std::logic_error("Could not find CurlGet when calling CurlStack::remove");
if (curl_multi_add_handle((CURLM*)m_handle, (*itr)->handle()) > 0)
throw torrent::internal_error("Error calling curl_multi_add_handle.");
m_size--;
m_getList.erase(itr);
} else {
m_active--;
}
}
void
@@ -146,9 +173,4 @@ CurlStack::global_cleanup() {
curl_global_cleanup();
}
CurlStack::SlotFactory
CurlStack::get_http_factory() {
return sigc::bind(sigc::ptr_fun(&CurlGet::new_object), this);
}
}
+38 -13
View File
@@ -37,7 +37,7 @@
#ifndef RTORRENT_CORE_CURL_STACK_H
#define RTORRENT_CORE_CURL_STACK_H
#include <list>
#include <deque>
#include <string>
#include <sigc++/slot.h>
@@ -45,31 +45,56 @@ namespace core {
class CurlGet;
class CurlStack {
// 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
// deque.
//
// This should fit well with the use-case of a http stack, thus
// we get most of the cache locality benefits of a vector with fast
// removal of elements.
class CurlStack : std::deque<CurlGet*> {
public:
friend class CurlGet;
typedef std::list<CurlGet*> CurlGetList;
typedef sigc::slot0<CurlGet*> SlotFactory;
typedef std::deque<CurlGet*> base_type;
using base_type::value_type;
using base_type::iterator;
using base_type::const_iterator;
using base_type::reverse_iterator;
using base_type::const_reverse_iterator;
using base_type::begin;
using base_type::end;
using base_type::rbegin;
using base_type::rend;
using base_type::back;
using base_type::front;
using base_type::size;
using base_type::empty;
CurlStack();
~CurlStack();
int get_size() const { return m_size; }
bool is_busy() const { return !m_getList.empty(); }
CurlGet* new_object();
void perform();
// TODO: Set fd_set's only once?
unsigned int fdset(fd_set* readfds, fd_set* writefds, fd_set* exceptfds);
SlotFactory get_http_factory();
unsigned int active() const { return m_active; }
unsigned int max_active() const { return m_maxActive; }
void set_max_active(unsigned int a) { m_maxActive = a; }
const std::string& user_agent() const { return m_userAgent; }
void set_user_agent(const std::string& s) { m_userAgent = s; }
const std::string& user_agent() const { return m_userAgent; }
void set_user_agent(const std::string& s) { m_userAgent = s; }
const std::string& http_proxy() const { return m_httpProxy; }
void set_http_proxy(const std::string& s) { m_httpProxy = s; }
const std::string& http_proxy() const { return m_httpProxy; }
void set_http_proxy(const std::string& s) { m_httpProxy = s; }
const std::string& bind_address() const { return m_bindAddress; }
void set_bind_address(const std::string& s) { m_bindAddress = s; }
@@ -87,8 +112,8 @@ class CurlStack {
void* m_handle;
int m_size;
CurlGetList m_getList;
unsigned int m_active;
unsigned int m_maxActive;
std::string m_userAgent;
std::string m_httpProxy;
-7
View File
@@ -51,12 +51,9 @@ namespace core {
HttpQueue::iterator
HttpQueue::insert(const std::string& url, std::iostream* s) {
std::auto_ptr<CurlGet> h(m_slotFactory());
//std::auto_ptr<std::stringstream> s(new std::stringstream);
h->set_url(url);
//h->set_stream(s.get());
h->set_stream(s);
h->set_user_agent("rtorrent/" VERSION);
iterator itr = Base::insert(end(), h.get());
@@ -66,8 +63,6 @@ HttpQueue::insert(const std::string& url, std::iostream* s) {
(*itr)->start();
h.release();
//s.release();
m_signalInsert.emit(*itr);
return itr;
@@ -77,9 +72,7 @@ void
HttpQueue::erase(iterator itr) {
m_signalErase.emit(*itr);
//delete (*itr)->get_stream();
delete *itr;
Base::erase(itr);
}
+2 -3
View File
@@ -36,7 +36,6 @@
#include "config.h"
#include <stdexcept>
#include <cstdio>
#include <cstring>
#include <fstream>
@@ -231,8 +230,8 @@ Manager::initialize_first() {
// Most of this should be possible to move out.
void
Manager::initialize_second() {
torrent::Http::set_factory(m_pollManager->get_http_stack()->get_http_factory());
m_httpQueue->slot_factory(m_pollManager->get_http_stack()->get_http_factory());
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));
CurlStack::global_init();
+1 -1
View File
@@ -67,7 +67,7 @@ PollManagerEPoll::poll(rak::timer timeout) {
torrent::perform();
timeout = std::min(timeout, rak::timer(torrent::next_timeout())) + 1000;
if (m_httpStack.is_busy()) {
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
+1 -1
View File
@@ -68,7 +68,7 @@ PollManagerKQueue::poll(rak::timer timeout) {
torrent::perform();
timeout = std::min(timeout, rak::timer(torrent::next_timeout())) + 1000;
if (m_httpStack.is_busy()) {
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
+2 -2
View File
@@ -77,7 +77,7 @@ PollManagerSelect::poll(rak::timer timeout) {
unsigned int maxFd = static_cast<torrent::PollSelect*>(m_poll)->fdset(m_readSet, m_writeSet, m_errorSet);
if (m_httpStack.is_busy())
if (!m_httpStack.empty())
maxFd = std::max(maxFd, m_httpStack.fdset(m_readSet, m_writeSet, m_errorSet));
timeval t = timeout.tval();
@@ -85,7 +85,7 @@ PollManagerSelect::poll(rak::timer timeout) {
if (select(maxFd + 1, m_readSet, m_writeSet, m_errorSet, &t) == -1)
return check_error();
if (m_httpStack.is_busy())
if (!m_httpStack.empty())
m_httpStack.perform();
torrent::perform();