Extract rpc_callbacks class
This commit is contained in:
parent
d5b3379def
commit
358f392992
6 changed files with 171 additions and 123 deletions
|
|
@ -135,6 +135,8 @@ add_library(
|
||||||
rep_tiers.cpp
|
rep_tiers.cpp
|
||||||
request_aggregator.hpp
|
request_aggregator.hpp
|
||||||
request_aggregator.cpp
|
request_aggregator.cpp
|
||||||
|
rpc_callbacks.hpp
|
||||||
|
rpc_callbacks.cpp
|
||||||
scheduler/bucket.cpp
|
scheduler/bucket.cpp
|
||||||
scheduler/bucket.hpp
|
scheduler/bucket.hpp
|
||||||
scheduler/component.hpp
|
scheduler/component.hpp
|
||||||
|
|
|
||||||
|
|
@ -33,6 +33,7 @@ class recently_cemented_cache;
|
||||||
class recently_confirmed_cache;
|
class recently_confirmed_cache;
|
||||||
class rep_crawler;
|
class rep_crawler;
|
||||||
class rep_tiers;
|
class rep_tiers;
|
||||||
|
class rpc_callbacks;
|
||||||
class telemetry;
|
class telemetry;
|
||||||
class unchecked_map;
|
class unchecked_map;
|
||||||
class stats;
|
class stats;
|
||||||
|
|
|
||||||
|
|
@ -30,6 +30,7 @@
|
||||||
#include <nano/node/peer_history.hpp>
|
#include <nano/node/peer_history.hpp>
|
||||||
#include <nano/node/portmapping.hpp>
|
#include <nano/node/portmapping.hpp>
|
||||||
#include <nano/node/request_aggregator.hpp>
|
#include <nano/node/request_aggregator.hpp>
|
||||||
|
#include <nano/node/rpc_callbacks.hpp>
|
||||||
#include <nano/node/scheduler/component.hpp>
|
#include <nano/node/scheduler/component.hpp>
|
||||||
#include <nano/node/scheduler/hinted.hpp>
|
#include <nano/node/scheduler/hinted.hpp>
|
||||||
#include <nano/node/scheduler/manual.hpp>
|
#include <nano/node/scheduler/manual.hpp>
|
||||||
|
|
@ -193,6 +194,8 @@ nano::node::node (std::shared_ptr<boost::asio::io_context> io_ctx_a, std::filesy
|
||||||
peer_history{ *peer_history_impl },
|
peer_history{ *peer_history_impl },
|
||||||
monitor_impl{ std::make_unique<nano::monitor> (config.monitor, *this) },
|
monitor_impl{ std::make_unique<nano::monitor> (config.monitor, *this) },
|
||||||
monitor{ *monitor_impl },
|
monitor{ *monitor_impl },
|
||||||
|
rpc_callbacks_impl{ std::make_unique<nano::rpc_callbacks> (*this) },
|
||||||
|
rpc_callbacks{ *rpc_callbacks_impl },
|
||||||
startup_time{ std::chrono::steady_clock::now () },
|
startup_time{ std::chrono::steady_clock::now () },
|
||||||
node_seq{ seq }
|
node_seq{ seq }
|
||||||
{
|
{
|
||||||
|
|
@ -250,66 +253,6 @@ nano::node::node (std::shared_ptr<boost::asio::io_context> io_ctx_a, std::filesy
|
||||||
network.disconnect_observer = [this] () {
|
network.disconnect_observer = [this] () {
|
||||||
observers.disconnect.notify ();
|
observers.disconnect.notify ();
|
||||||
};
|
};
|
||||||
if (!config.callback_address.empty ())
|
|
||||||
{
|
|
||||||
observers.blocks.add ([this] (nano::election_status const & status_a, std::vector<nano::vote_with_weight_info> const & votes_a, nano::account const & account_a, nano::amount const & amount_a, bool is_state_send_a, bool is_state_epoch_a) {
|
|
||||||
auto block_a (status_a.winner);
|
|
||||||
if ((status_a.type == nano::election_status_type::active_confirmed_quorum || status_a.type == nano::election_status_type::active_confirmation_height))
|
|
||||||
{
|
|
||||||
auto node_l (shared_from_this ());
|
|
||||||
io_ctx.post ([node_l, block_a, account_a, amount_a, is_state_send_a, is_state_epoch_a] () {
|
|
||||||
boost::property_tree::ptree event;
|
|
||||||
event.add ("account", account_a.to_account ());
|
|
||||||
event.add ("hash", block_a->hash ().to_string ());
|
|
||||||
std::string block_text;
|
|
||||||
block_a->serialize_json (block_text);
|
|
||||||
event.add ("block", block_text);
|
|
||||||
event.add ("amount", amount_a.to_string_dec ());
|
|
||||||
if (is_state_send_a)
|
|
||||||
{
|
|
||||||
event.add ("is_send", is_state_send_a);
|
|
||||||
event.add ("subtype", "send");
|
|
||||||
}
|
|
||||||
// Subtype field
|
|
||||||
else if (block_a->type () == nano::block_type::state)
|
|
||||||
{
|
|
||||||
if (block_a->is_change ())
|
|
||||||
{
|
|
||||||
event.add ("subtype", "change");
|
|
||||||
}
|
|
||||||
else if (is_state_epoch_a)
|
|
||||||
{
|
|
||||||
debug_assert (amount_a == 0 && node_l->ledger.is_epoch_link (block_a->link_field ().value ()));
|
|
||||||
event.add ("subtype", "epoch");
|
|
||||||
}
|
|
||||||
else
|
|
||||||
{
|
|
||||||
event.add ("subtype", "receive");
|
|
||||||
}
|
|
||||||
}
|
|
||||||
std::stringstream ostream;
|
|
||||||
boost::property_tree::write_json (ostream, event);
|
|
||||||
ostream.flush ();
|
|
||||||
auto body (std::make_shared<std::string> (ostream.str ()));
|
|
||||||
auto address (node_l->config.callback_address);
|
|
||||||
auto port (node_l->config.callback_port);
|
|
||||||
auto target (std::make_shared<std::string> (node_l->config.callback_target));
|
|
||||||
auto resolver (std::make_shared<boost::asio::ip::tcp::resolver> (node_l->io_ctx));
|
|
||||||
resolver->async_resolve (boost::asio::ip::tcp::resolver::query (address, std::to_string (port)), [node_l, address, port, target, body, resolver] (boost::system::error_code const & ec, boost::asio::ip::tcp::resolver::iterator i_a) {
|
|
||||||
if (!ec)
|
|
||||||
{
|
|
||||||
node_l->do_rpc_callback (i_a, address, port, target, body, resolver);
|
|
||||||
}
|
|
||||||
else
|
|
||||||
{
|
|
||||||
node_l->logger.error (nano::log::type::rpc_callbacks, "Error resolving callback: {}:{} ({})", address, port, ec.message ());
|
|
||||||
node_l->stats.inc (nano::stat::type::error, nano::stat::detail::http_callback, nano::stat::dir::out);
|
|
||||||
}
|
|
||||||
});
|
|
||||||
});
|
|
||||||
}
|
|
||||||
});
|
|
||||||
}
|
|
||||||
|
|
||||||
observers.channel_connected.add ([this] (std::shared_ptr<nano::transport::channel> const & channel) {
|
observers.channel_connected.add ([this] (std::shared_ptr<nano::transport::channel> const & channel) {
|
||||||
network.send_keepalive_self (channel);
|
network.send_keepalive_self (channel);
|
||||||
|
|
@ -472,68 +415,6 @@ nano::node::~node ()
|
||||||
stop ();
|
stop ();
|
||||||
}
|
}
|
||||||
|
|
||||||
// TODO: Move to a separate class
|
|
||||||
void nano::node::do_rpc_callback (boost::asio::ip::tcp::resolver::iterator i_a, std::string const & address, uint16_t port, std::shared_ptr<std::string> const & target, std::shared_ptr<std::string> const & body, std::shared_ptr<boost::asio::ip::tcp::resolver> const & resolver)
|
|
||||||
{
|
|
||||||
if (i_a != boost::asio::ip::tcp::resolver::iterator{})
|
|
||||||
{
|
|
||||||
auto node_l (shared_from_this ());
|
|
||||||
auto sock (std::make_shared<boost::asio::ip::tcp::socket> (node_l->io_ctx));
|
|
||||||
sock->async_connect (i_a->endpoint (), [node_l, target, body, sock, address, port, i_a, resolver] (boost::system::error_code const & ec) mutable {
|
|
||||||
if (!ec)
|
|
||||||
{
|
|
||||||
auto req (std::make_shared<boost::beast::http::request<boost::beast::http::string_body>> ());
|
|
||||||
req->method (boost::beast::http::verb::post);
|
|
||||||
req->target (*target);
|
|
||||||
req->version (11);
|
|
||||||
req->insert (boost::beast::http::field::host, address);
|
|
||||||
req->insert (boost::beast::http::field::content_type, "application/json");
|
|
||||||
req->body () = *body;
|
|
||||||
req->prepare_payload ();
|
|
||||||
boost::beast::http::async_write (*sock, *req, [node_l, sock, address, port, req, i_a, target, body, resolver] (boost::system::error_code const & ec, std::size_t bytes_transferred) mutable {
|
|
||||||
if (!ec)
|
|
||||||
{
|
|
||||||
auto sb (std::make_shared<boost::beast::flat_buffer> ());
|
|
||||||
auto resp (std::make_shared<boost::beast::http::response<boost::beast::http::string_body>> ());
|
|
||||||
boost::beast::http::async_read (*sock, *sb, *resp, [node_l, sb, resp, sock, address, port, i_a, target, body, resolver] (boost::system::error_code const & ec, std::size_t bytes_transferred) mutable {
|
|
||||||
if (!ec)
|
|
||||||
{
|
|
||||||
if (boost::beast::http::to_status_class (resp->result ()) == boost::beast::http::status_class::successful)
|
|
||||||
{
|
|
||||||
node_l->stats.inc (nano::stat::type::http_callback, nano::stat::detail::initiate, nano::stat::dir::out);
|
|
||||||
}
|
|
||||||
else
|
|
||||||
{
|
|
||||||
node_l->logger.error (nano::log::type::rpc_callbacks, "Callback to {}:{} failed [status: {}]", address, port, nano::util::to_str (resp->result ()));
|
|
||||||
node_l->stats.inc (nano::stat::type::error, nano::stat::detail::http_callback, nano::stat::dir::out);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
else
|
|
||||||
{
|
|
||||||
node_l->logger.error (nano::log::type::rpc_callbacks, "Unable to complete callback: {}:{} ({})", address, port, ec.message ());
|
|
||||||
node_l->stats.inc (nano::stat::type::error, nano::stat::detail::http_callback, nano::stat::dir::out);
|
|
||||||
};
|
|
||||||
});
|
|
||||||
}
|
|
||||||
else
|
|
||||||
{
|
|
||||||
node_l->logger.error (nano::log::type::rpc_callbacks, "Unable to send callback: {}:{} ({})", address, port, ec.message ());
|
|
||||||
node_l->stats.inc (nano::stat::type::error, nano::stat::detail::http_callback, nano::stat::dir::out);
|
|
||||||
}
|
|
||||||
});
|
|
||||||
}
|
|
||||||
else
|
|
||||||
{
|
|
||||||
node_l->logger.error (nano::log::type::rpc_callbacks, "Unable to connect to callback address({}): {}:{} ({})", address, i_a->endpoint ().address ().to_string (), port, ec.message ());
|
|
||||||
node_l->stats.inc (nano::stat::type::error, nano::stat::detail::http_callback, nano::stat::dir::out);
|
|
||||||
++i_a;
|
|
||||||
|
|
||||||
node_l->do_rpc_callback (i_a, address, port, target, body, resolver);
|
|
||||||
}
|
|
||||||
});
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
bool nano::node::copy_with_compaction (std::filesystem::path const & destination)
|
bool nano::node::copy_with_compaction (std::filesystem::path const & destination)
|
||||||
{
|
{
|
||||||
return store.copy_db (destination);
|
return store.copy_db (destination);
|
||||||
|
|
|
||||||
|
|
@ -82,7 +82,6 @@ public:
|
||||||
bool block_confirmed_or_being_confirmed (nano::secure::transaction const &, nano::block_hash const &);
|
bool block_confirmed_or_being_confirmed (nano::secure::transaction const &, nano::block_hash const &);
|
||||||
bool block_confirmed_or_being_confirmed (nano::block_hash const &);
|
bool block_confirmed_or_being_confirmed (nano::block_hash const &);
|
||||||
|
|
||||||
void do_rpc_callback (boost::asio::ip::tcp::resolver::iterator i_a, std::string const &, uint16_t, std::shared_ptr<std::string> const &, std::shared_ptr<std::string> const &, std::shared_ptr<boost::asio::ip::tcp::resolver> const &);
|
|
||||||
bool online () const;
|
bool online () const;
|
||||||
bool init_error () const;
|
bool init_error () const;
|
||||||
std::pair<uint64_t, std::unordered_map<nano::account, nano::uint128_t>> get_bootstrap_weights () const;
|
std::pair<uint64_t, std::unordered_map<nano::account, nano::uint128_t>> get_bootstrap_weights () const;
|
||||||
|
|
@ -201,6 +200,8 @@ public:
|
||||||
nano::peer_history & peer_history;
|
nano::peer_history & peer_history;
|
||||||
std::unique_ptr<nano::monitor> monitor_impl;
|
std::unique_ptr<nano::monitor> monitor_impl;
|
||||||
nano::monitor & monitor;
|
nano::monitor & monitor;
|
||||||
|
std::unique_ptr<nano::rpc_callbacks> rpc_callbacks_impl;
|
||||||
|
nano::rpc_callbacks & rpc_callbacks;
|
||||||
|
|
||||||
public:
|
public:
|
||||||
std::chrono::steady_clock::time_point const startup_time;
|
std::chrono::steady_clock::time_point const startup_time;
|
||||||
|
|
|
||||||
139
nano/node/rpc_callbacks.cpp
Normal file
139
nano/node/rpc_callbacks.cpp
Normal file
|
|
@ -0,0 +1,139 @@
|
||||||
|
#include <nano/lib/block_type.hpp>
|
||||||
|
#include <nano/node/node.hpp>
|
||||||
|
#include <nano/node/rpc_callbacks.hpp>
|
||||||
|
#include <nano/secure/ledger.hpp>
|
||||||
|
|
||||||
|
nano::rpc_callbacks::rpc_callbacks (nano::node & node_a) :
|
||||||
|
node{ node_a },
|
||||||
|
config{ node_a.config },
|
||||||
|
observers{ node_a.observers },
|
||||||
|
ledger{ node_a.ledger },
|
||||||
|
logger{ node_a.logger },
|
||||||
|
stats{ node_a.stats }
|
||||||
|
{
|
||||||
|
if (!config.callback_address.empty ())
|
||||||
|
{
|
||||||
|
logger.info (nano::log::type::rpc_callbacks, "RPC callbacks enabled on {}:{}", config.callback_address, config.callback_port);
|
||||||
|
setup_callbacks ();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
void nano::rpc_callbacks::setup_callbacks ()
|
||||||
|
{
|
||||||
|
observers.blocks.add ([this] (nano::election_status const & status_a, std::vector<nano::vote_with_weight_info> const & votes_a, nano::account const & account_a, nano::amount const & amount_a, bool is_state_send_a, bool is_state_epoch_a) {
|
||||||
|
auto block_a (status_a.winner);
|
||||||
|
if ((status_a.type == nano::election_status_type::active_confirmed_quorum || status_a.type == nano::election_status_type::active_confirmation_height))
|
||||||
|
{
|
||||||
|
node.workers.post ([this, block_a, account_a, amount_a, is_state_send_a, is_state_epoch_a] () {
|
||||||
|
boost::property_tree::ptree event;
|
||||||
|
event.add ("account", account_a.to_account ());
|
||||||
|
event.add ("hash", block_a->hash ().to_string ());
|
||||||
|
std::string block_text;
|
||||||
|
block_a->serialize_json (block_text);
|
||||||
|
event.add ("block", block_text);
|
||||||
|
event.add ("amount", amount_a.to_string_dec ());
|
||||||
|
if (is_state_send_a)
|
||||||
|
{
|
||||||
|
event.add ("is_send", is_state_send_a);
|
||||||
|
event.add ("subtype", "send");
|
||||||
|
}
|
||||||
|
// Subtype field
|
||||||
|
else if (block_a->type () == nano::block_type::state)
|
||||||
|
{
|
||||||
|
if (block_a->is_change ())
|
||||||
|
{
|
||||||
|
event.add ("subtype", "change");
|
||||||
|
}
|
||||||
|
else if (is_state_epoch_a)
|
||||||
|
{
|
||||||
|
debug_assert (amount_a == 0 && ledger.is_epoch_link (block_a->link_field ().value ()));
|
||||||
|
event.add ("subtype", "epoch");
|
||||||
|
}
|
||||||
|
else
|
||||||
|
{
|
||||||
|
event.add ("subtype", "receive");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
std::stringstream ostream;
|
||||||
|
boost::property_tree::write_json (ostream, event);
|
||||||
|
ostream.flush ();
|
||||||
|
auto body (std::make_shared<std::string> (ostream.str ()));
|
||||||
|
auto address (config.callback_address);
|
||||||
|
auto port (config.callback_port);
|
||||||
|
auto target (std::make_shared<std::string> (config.callback_target));
|
||||||
|
auto resolver (std::make_shared<boost::asio::ip::tcp::resolver> (node.io_ctx));
|
||||||
|
resolver->async_resolve (boost::asio::ip::tcp::resolver::query (address, std::to_string (port)), [this, address, port, target, body, resolver] (boost::system::error_code const & ec, boost::asio::ip::tcp::resolver::iterator i_a) {
|
||||||
|
if (!ec)
|
||||||
|
{
|
||||||
|
do_rpc_callback (i_a, address, port, target, body, resolver);
|
||||||
|
}
|
||||||
|
else
|
||||||
|
{
|
||||||
|
logger.error (nano::log::type::rpc_callbacks, "Error resolving callback: {}:{} ({})", address, port, ec.message ());
|
||||||
|
stats.inc (nano::stat::type::error, nano::stat::detail::http_callback, nano::stat::dir::out);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
});
|
||||||
|
}
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
void nano::rpc_callbacks::do_rpc_callback (boost::asio::ip::tcp::resolver::iterator i_a, std::string const & address, uint16_t port, std::shared_ptr<std::string> const & target, std::shared_ptr<std::string> const & body, std::shared_ptr<boost::asio::ip::tcp::resolver> const & resolver)
|
||||||
|
{
|
||||||
|
if (i_a != boost::asio::ip::tcp::resolver::iterator{})
|
||||||
|
{
|
||||||
|
auto sock (std::make_shared<boost::asio::ip::tcp::socket> (node.io_ctx));
|
||||||
|
sock->async_connect (i_a->endpoint (), [this, target, body, sock, address, port, i_a, resolver] (boost::system::error_code const & ec) mutable {
|
||||||
|
if (!ec)
|
||||||
|
{
|
||||||
|
auto req (std::make_shared<boost::beast::http::request<boost::beast::http::string_body>> ());
|
||||||
|
req->method (boost::beast::http::verb::post);
|
||||||
|
req->target (*target);
|
||||||
|
req->version (11);
|
||||||
|
req->insert (boost::beast::http::field::host, address);
|
||||||
|
req->insert (boost::beast::http::field::content_type, "application/json");
|
||||||
|
req->body () = *body;
|
||||||
|
req->prepare_payload ();
|
||||||
|
boost::beast::http::async_write (*sock, *req, [this, sock, address, port, req, i_a, target, body, resolver] (boost::system::error_code const & ec, std::size_t bytes_transferred) mutable {
|
||||||
|
if (!ec)
|
||||||
|
{
|
||||||
|
auto sb (std::make_shared<boost::beast::flat_buffer> ());
|
||||||
|
auto resp (std::make_shared<boost::beast::http::response<boost::beast::http::string_body>> ());
|
||||||
|
boost::beast::http::async_read (*sock, *sb, *resp, [this, sb, resp, sock, address, port, i_a, target, body, resolver] (boost::system::error_code const & ec, std::size_t bytes_transferred) mutable {
|
||||||
|
if (!ec)
|
||||||
|
{
|
||||||
|
if (boost::beast::http::to_status_class (resp->result ()) == boost::beast::http::status_class::successful)
|
||||||
|
{
|
||||||
|
stats.inc (nano::stat::type::http_callback, nano::stat::detail::initiate, nano::stat::dir::out);
|
||||||
|
}
|
||||||
|
else
|
||||||
|
{
|
||||||
|
logger.error (nano::log::type::rpc_callbacks, "Callback to {}:{} failed [status: {}]", address, port, nano::util::to_str (resp->result ()));
|
||||||
|
stats.inc (nano::stat::type::error, nano::stat::detail::http_callback, nano::stat::dir::out);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
else
|
||||||
|
{
|
||||||
|
logger.error (nano::log::type::rpc_callbacks, "Unable to complete callback: {}:{} ({})", address, port, ec.message ());
|
||||||
|
stats.inc (nano::stat::type::error, nano::stat::detail::http_callback, nano::stat::dir::out);
|
||||||
|
};
|
||||||
|
});
|
||||||
|
}
|
||||||
|
else
|
||||||
|
{
|
||||||
|
logger.error (nano::log::type::rpc_callbacks, "Unable to send callback: {}:{} ({})", address, port, ec.message ());
|
||||||
|
stats.inc (nano::stat::type::error, nano::stat::detail::http_callback, nano::stat::dir::out);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
}
|
||||||
|
else
|
||||||
|
{
|
||||||
|
logger.error (nano::log::type::rpc_callbacks, "Unable to connect to callback address({}): {}:{} ({})", address, i_a->endpoint ().address ().to_string (), port, ec.message ());
|
||||||
|
stats.inc (nano::stat::type::error, nano::stat::detail::http_callback, nano::stat::dir::out);
|
||||||
|
++i_a;
|
||||||
|
|
||||||
|
do_rpc_callback (i_a, address, port, target, body, resolver);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
24
nano/node/rpc_callbacks.hpp
Normal file
24
nano/node/rpc_callbacks.hpp
Normal file
|
|
@ -0,0 +1,24 @@
|
||||||
|
#pragma once
|
||||||
|
|
||||||
|
#include <nano/node/fwd.hpp>
|
||||||
|
|
||||||
|
namespace nano
|
||||||
|
{
|
||||||
|
class rpc_callbacks
|
||||||
|
{
|
||||||
|
public:
|
||||||
|
explicit rpc_callbacks (nano::node &);
|
||||||
|
|
||||||
|
private: // Dependencies
|
||||||
|
nano::node_config const & config;
|
||||||
|
nano::node & node;
|
||||||
|
nano::node_observers & observers;
|
||||||
|
nano::ledger & ledger;
|
||||||
|
nano::logger & logger;
|
||||||
|
nano::stats & stats;
|
||||||
|
|
||||||
|
private:
|
||||||
|
void setup_callbacks ();
|
||||||
|
void do_rpc_callback (boost::asio::ip::tcp::resolver::iterator i_a, std::string const &, uint16_t, std::shared_ptr<std::string> const &, std::shared_ptr<std::string> const &, std::shared_ptr<boost::asio::ip::tcp::resolver> const &);
|
||||||
|
};
|
||||||
|
}
|
||||||
Loading…
Add table
Add a link
Reference in a new issue