Throttle all commands sent to IRC servers

fix #3354
This commit is contained in:
louiz’
2018-06-25 22:54:32 +02:00
parent ba97c442a8
commit 09b10cc801
9 changed files with 172 additions and 34 deletions
+4 -1
View File
@@ -86,13 +86,16 @@ class Database
struct Address: Column<std::string> { static constexpr auto name = "address_"; };
struct ThrottleLimit: Column<std::size_t> { static constexpr auto name = "throttlelimit_";
ThrottleLimit(): Column<std::size_t>(10) {} };
using MucLogLineTable = Table<Id, Uuid, Owner, IrcChanName, IrcServerName, Date, Body, Nick>;
using MucLogLine = MucLogLineTable::RowType;
using GlobalOptionsTable = Table<Id, Owner, MaxHistoryLength, RecordHistory, GlobalPersistent>;
using GlobalOptions = GlobalOptionsTable::RowType;
using IrcServerOptionsTable = Table<Id, Owner, Server, Pass, TlsPorts, Ports, Username, Realname, VerifyCert, TrustedFingerprint, EncodingOut, EncodingIn, MaxHistoryLength, Address, Nick>;
using IrcServerOptionsTable = Table<Id, Owner, Server, Pass, TlsPorts, Ports, Username, Realname, VerifyCert, TrustedFingerprint, EncodingOut, EncodingIn, MaxHistoryLength, Address, Nick, ThrottleLimit>;
using IrcServerOptions = IrcServerOptionsTable::RowType;
using IrcChannelOptionsTable = Table<Id, Owner, Server, Channel, EncodingOut, EncodingIn, MaxHistoryLength, Persistent, RecordHistoryOptional>;
+42 -21
View File
@@ -135,7 +135,7 @@ IrcClient::IrcClient(std::shared_ptr<Poller>& poller, std::string hostname,
std::string realname, std::string user_hostname,
Bridge& bridge):
TCPClientSocketHandler(poller),
hostname(std::move(hostname)),
hostname(hostname),
user_hostname(std::move(user_hostname)),
username(std::move(username)),
realname(std::move(realname)),
@@ -143,7 +143,14 @@ IrcClient::IrcClient(std::shared_ptr<Poller>& poller, std::string hostname,
bridge(bridge),
welcomed(false),
chanmodes({"", "", "", ""}),
chantypes({'#', '&'})
chantypes({'#', '&'}),
tokens_bucket(Database::get_irc_server_options(bridge.get_bare_jid(), hostname).col<Database::ThrottleLimit>(), 1s, [this]() {
if (message_queue.empty())
return true;
this->actual_send(std::move(this->message_queue.front()));
this->message_queue.pop_front();
return false;
}, "TokensBucket" + this->hostname + this->bridge.get_jid())
{
#ifdef USE_DATABASE
auto options = Database::get_irc_server_options(this->bridge.get_bare_jid(),
@@ -171,6 +178,7 @@ IrcClient::~IrcClient()
// This event may or may not exist (if we never got connected, it
// doesn't), but it's ok
TimedEventsManager::instance().cancel("PING" + this->hostname + this->bridge.get_jid());
TimedEventsManager::instance().cancel("TokensBucket" + this->hostname + this->bridge.get_jid());
}
void IrcClient::start()
@@ -390,25 +398,33 @@ void IrcClient::parse_in_buffer(const size_t)
}
}
void IrcClient::send_message(IrcMessage&& message)
void IrcClient::actual_send(const IrcMessage& message)
{
log_debug("IRC SENDING: (", this->get_hostname(), ") ", message);
std::string res;
if (!message.prefix.empty())
res += ":" + std::move(message.prefix) + " ";
res += message.command;
for (const std::string& arg: message.arguments)
{
if (arg.find(' ') != std::string::npos ||
(!arg.empty() && arg[0] == ':'))
{
res += " :" + arg;
break;
}
res += " " + arg;
}
res += "\r\n";
this->send_data(std::move(res));
log_debug("IRC SENDING: (", this->get_hostname(), ") ", message);
std::string res;
if (!message.prefix.empty())
res += ":" + message.prefix + " ";
res += message.command;
for (const std::string& arg: message.arguments)
{
if (arg.find(' ') != std::string::npos
|| (!arg.empty() && arg[0] == ':'))
{
res += " :" + arg;
break;
}
res += " " + arg;
}
res += "\r\n";
this->send_data(std::move(res));
}
void IrcClient::send_message(IrcMessage message, bool throttle)
{
if (this->tokens_bucket.use_token() || !throttle)
this->actual_send(message);
else
message_queue.push_back(std::move(message));
}
void IrcClient::send_raw(const std::string& txt)
@@ -459,7 +475,7 @@ void IrcClient::send_topic_command(const std::string& chan_name, const std::stri
void IrcClient::send_quit_command(const std::string& reason)
{
this->send_message(IrcMessage("QUIT", {reason}));
this->send_message(IrcMessage("QUIT", {reason}), false);
}
void IrcClient::send_join_command(const std::string& chan_name, const std::string& password)
@@ -1225,6 +1241,11 @@ void IrcClient::on_channel_mode(const IrcMessage& message)
}
}
void IrcClient::set_throttle_limit(std::size_t limit)
{
this->tokens_bucket.set_limit(limit);
}
void IrcClient::on_user_mode(const IrcMessage& message)
{
this->bridge.send_xmpp_message(this->hostname, "",
+10 -2
View File
@@ -16,8 +16,10 @@
#include <vector>
#include <string>
#include <stack>
#include <deque>
#include <map>
#include <set>
#include <utils/tokens_bucket.hpp>
class Bridge;
@@ -84,8 +86,9 @@ public:
* (actually, into our out_buf and signal the poller that we want to wach
* for send events to be ready)
*/
void send_message(IrcMessage&& message);
void send_message(IrcMessage message, bool throttle=true);
void send_raw(const std::string& txt);
void actual_send(const IrcMessage& message);
/**
* Send the PONG irc command
*/
@@ -293,7 +296,7 @@ public:
const std::vector<char>& get_sorted_user_modes() const { return this->sorted_user_modes; }
std::set<char> get_chantypes() const { return this->chantypes; }
void set_throttle_limit(std::size_t limit);
/**
* Store the history limit that the client asked when joining this room.
*/
@@ -330,6 +333,10 @@ private:
* To communicate back with the bridge
*/
Bridge& bridge;
/**
* Where messaged are stored when they are throttled.
*/
std::deque<IrcMessage> message_queue{};
/**
* The list of joined channels, indexed by name
*/
@@ -389,6 +396,7 @@ private:
* the WebIRC protocole.
*/
Resolver dns_resolver;
TokensBucket tokens_bucket;
};
+2 -2
View File
@@ -14,9 +14,9 @@ public:
~IrcMessage() = default;
IrcMessage(const IrcMessage&) = delete;
IrcMessage(IrcMessage&&) = delete;
IrcMessage(IrcMessage&&) = default;
IrcMessage& operator=(const IrcMessage&) = delete;
IrcMessage& operator=(IrcMessage&&) = delete;
IrcMessage& operator=(IrcMessage&&) = default;
std::string prefix;
std::string command;
+58
View File
@@ -0,0 +1,58 @@
/**
* Implementation of the token bucket algorithm.
*
* It uses a repetitive TimedEvent, started at construction, to fill the
* bucket.
*
* Every n seconds, it executes the given callback. If the callback
* returns true, we add a token (if the limit is not yet reached).
*
*/
#pragma once
#include <utils/timed_events.hpp>
#include <logger/logger.hpp>
class TokensBucket
{
public:
TokensBucket(std::size_t max_size, std::chrono::milliseconds fill_duration, std::function<bool()> callback, std::string name):
limit(max_size),
tokens(limit),
fill_duration(fill_duration),
callback(std::move(callback))
{
log_debug("creating TokensBucket with max size: ", max_size);
TimedEvent event(std::move(fill_duration), [this]() { this->add_token(); }, std::move(name));
TimedEventsManager::instance().add_event(std::move(event));
}
bool use_token()
{
if (this->tokens > 0)
{
this->tokens--;
return true;
}
else
return false;
}
void set_limit(std::size_t limit)
{
this->limit = limit;
}
private:
std::size_t limit;
std::size_t tokens;
std::chrono::milliseconds fill_duration;
std::function<bool()> callback;
void add_token()
{
if (this->callback() && this->tokens != limit)
this->tokens++;
}
};
+25 -1
View File
@@ -365,6 +365,15 @@ void ConfigureIrcServerStep1(XmppComponent&, AdhocSession& session, XmlNode& com
}
}
{
XmlSubNode throttle_limit(x, "field");
throttle_limit["var"] = "throttle_limit";
throttle_limit["type"] = "text-single";
throttle_limit["label"] = "Throttle limit";
XmlSubNode value(throttle_limit, "value");
value.set_inner(std::to_string(options.col<Database::ThrottleLimit>()));
}
{
XmlSubNode encoding_out(x, "field");
encoding_out["var"] = "encoding_out";
@@ -392,8 +401,10 @@ void ConfigureIrcServerStep1(XmppComponent&, AdhocSession& session, XmlNode& com
}
}
void ConfigureIrcServerStep2(XmppComponent&, AdhocSession& session, XmlNode& command_node)
void ConfigureIrcServerStep2(XmppComponent& xmpp_component, AdhocSession& session, XmlNode& command_node)
{
auto& biboumi_component = dynamic_cast<BiboumiComponent&>(xmpp_component);
const XmlNode* x = command_node.get_child("x", "jabber:x:data");
if (x)
{
@@ -474,6 +485,19 @@ void ConfigureIrcServerStep2(XmppComponent&, AdhocSession& session, XmlNode& com
else if (field->get_tag("var") == "realname" && value)
options.col<Database::Realname>() = value->get_inner();
else if (field->get_tag("var") == "throttle_limit" && value)
{
options.col<Database::ThrottleLimit>() = std::stoull(value->get_inner());
Bridge* bridge = biboumi_component.find_user_bridge(session.get_owner_jid());
if (bridge)
{
IrcClient* client = bridge->find_irc_client(server_domain);
if (client)
client->set_throttle_limit(options.col<Database::ThrottleLimit>());
}
}
else if (field->get_tag("var") == "encoding_out" && value)
options.col<Database::EncodingOut>() = value->get_inner();