352 lines
12 KiB
C++
352 lines
12 KiB
C++
#include <graphene/peerplays_sidechain/common/rpc_connection.hpp>
|
|
|
|
#include <regex>
|
|
#include <sstream>
|
|
|
|
#include <boost/asio/buffers_iterator.hpp>
|
|
#include <boost/asio/connect.hpp>
|
|
#include <boost/asio/ssl/error.hpp>
|
|
#include <boost/asio/ssl/stream.hpp>
|
|
#include <boost/beast/http.hpp>
|
|
#include <boost/property_tree/json_parser.hpp>
|
|
#include <boost/property_tree/ptree.hpp>
|
|
#include <boost/xpressive/xpressive.hpp>
|
|
|
|
#include <fc/log/logger.hpp>
|
|
|
|
#include <graphene/peerplays_sidechain/common/utils.hpp>
|
|
|
|
namespace graphene { namespace peerplays_sidechain {
|
|
|
|
struct rpc_reply {
|
|
uint16_t status;
|
|
std::string body;
|
|
};
|
|
|
|
class rpc_connection {
|
|
public:
|
|
rpc_connection(std::string _url, std::string _user, std::string _password, bool _debug_rpc_calls);
|
|
|
|
protected:
|
|
std::string retrieve_array_value_from_reply(std::string reply_str, std::string array_path, uint32_t idx);
|
|
std::string retrieve_value_from_reply(std::string reply_str, std::string value_path);
|
|
std::string send_post_request(std::string method, std::string params, bool show_log);
|
|
|
|
std::string url;
|
|
std::string user;
|
|
std::string password;
|
|
bool debug_rpc_calls;
|
|
|
|
std::string protocol;
|
|
std::string host;
|
|
std::string port;
|
|
std::string target;
|
|
std::string authorization;
|
|
|
|
uint32_t request_id;
|
|
|
|
private:
|
|
rpc_reply send_post_request(std::string body, bool show_log);
|
|
|
|
boost::beast::net::io_context ioc;
|
|
boost::beast::net::ip::tcp::resolver resolver;
|
|
boost::asio::ip::basic_resolver_results<boost::asio::ip::tcp> results;
|
|
};
|
|
|
|
rpc_connection::rpc_connection(std::string _url, std::string _user, std::string _password, bool _debug_rpc_calls) :
|
|
url(_url),
|
|
user(_user),
|
|
password(_password),
|
|
debug_rpc_calls(_debug_rpc_calls),
|
|
request_id(0),
|
|
resolver(ioc) {
|
|
|
|
std::string reg_expr = "^((?P<Protocol>https|http):\\/\\/)?(?P<Host>[a-zA-Z0-9\\-\\.]+)(:(?P<Port>\\d{1,5}))?(?P<Target>\\/.+)?";
|
|
boost::xpressive::sregex sr = boost::xpressive::sregex::compile(reg_expr);
|
|
|
|
boost::xpressive::smatch sm;
|
|
|
|
if (boost::xpressive::regex_search(url, sm, sr)) {
|
|
protocol = sm["Protocol"];
|
|
if (protocol.empty()) {
|
|
protocol = "http";
|
|
}
|
|
|
|
host = sm["Host"];
|
|
if (host.empty()) {
|
|
host + "localhost";
|
|
}
|
|
|
|
port = sm["Port"];
|
|
if (port.empty()) {
|
|
port = "80";
|
|
}
|
|
|
|
target = sm["Target"];
|
|
if (target.empty()) {
|
|
target = "/";
|
|
}
|
|
|
|
authorization = "Basic " + base64_encode(user + ":" + password);
|
|
|
|
results = resolver.resolve(host, port);
|
|
|
|
} else {
|
|
elog("Invalid URL: ${url}", ("url", url));
|
|
}
|
|
}
|
|
|
|
std::string rpc_connection::retrieve_array_value_from_reply(std::string reply_str, std::string array_path, uint32_t idx) {
|
|
if (reply_str.empty()) {
|
|
wlog("RPC call ${function}, empty reply string", ("function", __FUNCTION__));
|
|
return "";
|
|
}
|
|
|
|
try {
|
|
std::stringstream ss(reply_str);
|
|
boost::property_tree::ptree json;
|
|
boost::property_tree::read_json(ss, json);
|
|
if (json.find("result") == json.not_found()) {
|
|
return "";
|
|
}
|
|
|
|
auto json_result = json.get_child_optional("result");
|
|
if (json_result) {
|
|
boost::property_tree::ptree array_ptree = json_result.get();
|
|
if (!array_path.empty()) {
|
|
array_ptree = json_result.get().get_child(array_path);
|
|
}
|
|
uint32_t array_el_idx = -1;
|
|
for (const auto &array_el : array_ptree) {
|
|
array_el_idx = array_el_idx + 1;
|
|
if (array_el_idx == idx) {
|
|
std::stringstream ss_res;
|
|
boost::property_tree::json_parser::write_json(ss_res, array_el.second);
|
|
return ss_res.str();
|
|
}
|
|
}
|
|
}
|
|
} catch (const boost::property_tree::json_parser::json_parser_error &e) {
|
|
wlog("RPC call ${function} failed: ${e}", ("function", __FUNCTION__)("e", e.what()));
|
|
}
|
|
|
|
return "";
|
|
}
|
|
|
|
std::string rpc_connection::retrieve_value_from_reply(std::string reply_str, std::string value_path) {
|
|
if (reply_str.empty()) {
|
|
wlog("RPC call ${function}, empty reply string", ("function", __FUNCTION__));
|
|
return "";
|
|
}
|
|
|
|
try {
|
|
std::stringstream ss(reply_str);
|
|
boost::property_tree::ptree json;
|
|
boost::property_tree::read_json(ss, json);
|
|
if (json.find("result") == json.not_found()) {
|
|
return "";
|
|
}
|
|
|
|
auto json_result = json.get_child_optional("result");
|
|
if (json_result) {
|
|
return json_result.get().get<std::string>(value_path);
|
|
}
|
|
|
|
return json.get<std::string>("result");
|
|
} catch (const boost::property_tree::json_parser::json_parser_error &e) {
|
|
wlog("RPC call ${function} failed: ${e}", ("function", __FUNCTION__)("e", e.what()));
|
|
}
|
|
|
|
return "";
|
|
}
|
|
|
|
std::string rpc_connection::send_post_request(std::string method, std::string params, bool show_log) {
|
|
std::stringstream body;
|
|
|
|
request_id = request_id + 1;
|
|
|
|
body << "{ \"jsonrpc\": \"2.0\", \"id\": " << request_id << ", \"method\": \"" << method << "\"";
|
|
|
|
if (!params.empty()) {
|
|
body << ", \"params\": " << params;
|
|
}
|
|
|
|
body << " }";
|
|
|
|
try {
|
|
const auto reply = send_post_request(body.str(), show_log);
|
|
|
|
if (reply.body.empty()) {
|
|
wlog("RPC call ${function} failed", ("function", __FUNCTION__));
|
|
return "";
|
|
}
|
|
|
|
std::stringstream ss(std::string(reply.body.begin(), reply.body.end()));
|
|
boost::property_tree::ptree json;
|
|
boost::property_tree::read_json(ss, json);
|
|
|
|
if (json.count("error") && !json.get_child("error").empty()) {
|
|
wlog("RPC call ${function} with body ${body} failed with reply '${msg}'", ("function", __FUNCTION__)("body", body.str())("msg", ss.str()));
|
|
}
|
|
|
|
if (reply.status == 200) {
|
|
return ss.str();
|
|
}
|
|
} catch (const boost::system::system_error &e) {
|
|
elog("RPC call ${function} failed: ${e}", ("function", __FUNCTION__)("e", e.what()));
|
|
}
|
|
|
|
return "";
|
|
}
|
|
|
|
rpc_reply rpc_connection::send_post_request(std::string body, bool show_log) {
|
|
|
|
// These object is used as a context for ssl connection
|
|
boost::asio::ssl::context ctx(boost::asio::ssl::context::tlsv12_client);
|
|
|
|
boost::beast::net::ssl::stream<boost::beast::tcp_stream> ssl_tcp_stream(ioc, ctx);
|
|
boost::beast::tcp_stream tcp_stream(ioc);
|
|
|
|
// Set SNI Hostname (many hosts need this to handshake successfully)
|
|
if (protocol == "https") {
|
|
if (!SSL_set_tlsext_host_name(ssl_tcp_stream.native_handle(), host.c_str())) {
|
|
boost::beast::error_code ec{static_cast<int>(::ERR_get_error()), boost::asio::error::get_ssl_category()};
|
|
throw boost::beast::system_error{ec};
|
|
}
|
|
ctx.set_default_verify_paths();
|
|
ctx.set_verify_mode(boost::asio::ssl::verify_peer);
|
|
}
|
|
|
|
// Make the connection on the IP address we get from a lookup
|
|
if (protocol == "https") {
|
|
boost::beast::get_lowest_layer(ssl_tcp_stream).connect(results);
|
|
ssl_tcp_stream.handshake(boost::beast::net::ssl::stream_base::client);
|
|
} else {
|
|
tcp_stream.connect(results);
|
|
}
|
|
|
|
// Set up an HTTP GET request message
|
|
boost::beast::http::request<boost::beast::http::string_body> req{boost::beast::http::verb::post, target, 11};
|
|
req.set(boost::beast::http::field::host, host + ":" + port);
|
|
req.set(boost::beast::http::field::accept, "application/json");
|
|
req.set(boost::beast::http::field::authorization, authorization);
|
|
req.set(boost::beast::http::field::content_type, "application/json");
|
|
req.set(boost::beast::http::field::content_encoding, "utf-8");
|
|
req.set(boost::beast::http::field::content_length, body.length());
|
|
req.body() = body;
|
|
|
|
// Send the HTTP request to the remote host
|
|
if (protocol == "https")
|
|
boost::beast::http::write(ssl_tcp_stream, req);
|
|
else
|
|
boost::beast::http::write(tcp_stream, req);
|
|
|
|
// This buffer is used for reading and must be persisted
|
|
boost::beast::flat_buffer buffer;
|
|
|
|
// Declare a container to hold the response
|
|
boost::beast::http::response<boost::beast::http::dynamic_body> res;
|
|
|
|
// Receive the HTTP response
|
|
if (protocol == "https")
|
|
boost::beast::http::read(ssl_tcp_stream, buffer, res);
|
|
else
|
|
boost::beast::http::read(tcp_stream, buffer, res);
|
|
|
|
// Gracefully close the socket
|
|
boost::beast::error_code ec;
|
|
if (protocol == "https") {
|
|
boost::beast::get_lowest_layer(ssl_tcp_stream).close();
|
|
} else {
|
|
tcp_stream.socket().shutdown(boost::asio::ip::tcp::socket::shutdown_both, ec);
|
|
}
|
|
|
|
// not_connected happens sometimes. Also on ssl level some servers are managing
|
|
// connecntion close, so closing here will sometimes end up with error stream truncated
|
|
// so don't bother reporting it.
|
|
if (ec && ec != boost::beast::errc::not_connected && ec != boost::asio::ssl::error::stream_truncated)
|
|
throw boost::beast::system_error{ec};
|
|
|
|
std::string rbody{boost::asio::buffers_begin(res.body().data()),
|
|
boost::asio::buffers_end(res.body().data())};
|
|
rpc_reply reply;
|
|
reply.status = 200;
|
|
reply.body = rbody;
|
|
|
|
if (show_log) {
|
|
ilog("### Request URL: ${url}", ("url", url));
|
|
ilog("### Request: ${body}", ("body", body));
|
|
ilog("### Response: ${rbody}", ("rbody", rbody));
|
|
}
|
|
|
|
return reply;
|
|
}
|
|
|
|
rpc_client::rpc_client(const std::vector<std::string> &_urls, const std::vector<std::string> &_users, const std::vector<std::string> &_passwords, bool _debug_rpc_calls)
|
|
{
|
|
FC_ASSERT(_urls.size());
|
|
FC_ASSERT(_users.size() == _urls.size() && _passwords.size() == _urls.size());
|
|
for (size_t i=0; i < _urls.size(); i++)
|
|
connections.push_back(new rpc_connection(_urls[i], _users[i], _passwords[i], _debug_rpc_calls));
|
|
n_active_conn = 0;
|
|
}
|
|
|
|
void rpc_client::reselect_connection()
|
|
{
|
|
//ilog("n_active_rpc_client=${n}", ("n", n_active_rpc_client));
|
|
FC_ASSERT(connections.size());
|
|
|
|
int best_n = -1;
|
|
int best_quality = -1;
|
|
|
|
std::vector<uint64_t> head_block_numbers;
|
|
head_block_numbers.resize(rpc_clients.size());
|
|
|
|
for (size_t n=0; n < connections.size(); n++) {
|
|
rpc_connection *conn = connections[n];
|
|
int quality = 0;
|
|
head_block_numbers[n] = std::numeric_limits<uint64_t>::max();
|
|
|
|
// make the ping
|
|
fc::time_point t_sent = fc::time_point::now();
|
|
uint64_t head_block_number = ping(*conn);
|
|
fc::time_point t_received = fc::time_point::now();
|
|
if (head_block_number != std::numeric_limits<uint64_t>::max()) {
|
|
int t = (t_received - t_sent).count();
|
|
t += rand() % 10;
|
|
FC_ASSERT(t != -1);
|
|
head_block_numbers[n] = head_block_number;
|
|
static const int t_limit = 10*1000000; // 10 sec
|
|
if (t < t_limit)
|
|
quality = t_limit - t; // the less time, the higher quality
|
|
|
|
// look for the best quality
|
|
if (quality > best_quality) {
|
|
best_n = n;
|
|
best_quality = quality;
|
|
}
|
|
}
|
|
|
|
FC_ASSERT(best_n != -1 && best_quality != -1);
|
|
if (best_n != n_active_conn) { // if the best client is not the current one, ...
|
|
uint64_t active_head_block_number = head_block_numbers[n_active_conn];
|
|
if (active_head_block_number == std::numeric_limits<uint64_t>::max() // ... and the current one has no known head block...
|
|
|| head_block_numbers[best_n] >= active_head_block_number) { // ...or the best client's head is more recent than the current, ...
|
|
n_active_conn = best_n; // ...then select new one
|
|
ilog("!!! rpc connection reselected, now ${n}", ("n", n_active_conn));
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
rpc_connection &rpc_client::get_active_connection() const
|
|
{
|
|
return *connections[n_active_conn];
|
|
}
|
|
|
|
std::string rpc_client::send_post_request(std::string method, std::string params, bool show_log)
|
|
{
|
|
return get_active_connection().send_post_request(method, params, show_log);
|
|
}
|
|
|
|
}} // namespace graphene::peerplays_sidechain
|