1#include <xrpld/overlay/detail/PeerImp.h>
3#include <xrpld/app/consensus/RCLCxPeerPos.h>
4#include <xrpld/app/consensus/RCLValidations.h>
5#include <xrpld/app/ledger/InboundLedgers.h>
6#include <xrpld/app/ledger/InboundTransactions.h>
7#include <xrpld/app/ledger/LedgerMaster.h>
8#include <xrpld/app/ledger/LedgerNodeHelpers.h>
9#include <xrpld/app/ledger/TransactionMaster.h>
10#include <xrpld/app/misc/Transaction.h>
11#include <xrpld/app/misc/ValidatorList.h>
12#include <xrpld/overlay/Cluster.h>
13#include <xrpld/overlay/ClusterNode.h>
14#include <xrpld/overlay/Peer.h>
15#include <xrpld/overlay/ReduceRelayCommon.h>
16#include <xrpld/overlay/detail/Handshake.h>
17#include <xrpld/overlay/detail/OverlayImpl.h>
18#include <xrpld/overlay/detail/ProtocolMessage.h>
19#include <xrpld/overlay/detail/ProtocolVersion.h>
20#include <xrpld/overlay/detail/TrafficCount.h>
21#include <xrpld/overlay/detail/Tuning.h>
23#include <xrpl/basics/Blob.h>
24#include <xrpl/basics/Log.h>
25#include <xrpl/basics/SHAMapHash.h>
26#include <xrpl/basics/Slice.h>
27#include <xrpl/basics/ToString.h>
28#include <xrpl/basics/UptimeClock.h>
29#include <xrpl/basics/base64.h>
30#include <xrpl/basics/base_uint.h>
31#include <xrpl/basics/chrono.h>
32#include <xrpl/basics/random.h>
33#include <xrpl/basics/safe_cast.h>
34#include <xrpl/basics/strHex.h>
35#include <xrpl/beast/utility/Journal.h>
36#include <xrpl/beast/utility/Zero.h>
37#include <xrpl/beast/utility/instrumentation.h>
38#include <xrpl/consensus/Validations.h>
39#include <xrpl/core/HashRouter.h>
40#include <xrpl/core/Job.h>
41#include <xrpl/core/PerfLog.h>
42#include <xrpl/json/json_forwards.h>
43#include <xrpl/json/json_value.h>
44#include <xrpl/ledger/Ledger.h>
45#include <xrpl/peerfinder/Slot.h>
46#include <xrpl/peerfinder/Types.h>
47#include <xrpl/protocol/Feature.h>
48#include <xrpl/protocol/KeyType.h>
49#include <xrpl/protocol/LedgerHeader.h>
50#include <xrpl/protocol/Protocol.h>
51#include <xrpl/protocol/PublicKey.h>
52#include <xrpl/protocol/SField.h>
53#include <xrpl/protocol/STTx.h>
54#include <xrpl/protocol/Serializer.h>
55#include <xrpl/protocol/TxFlags.h>
56#include <xrpl/protocol/digest.h>
57#include <xrpl/protocol/jss.h>
58#include <xrpl/protocol/tokens.h>
59#include <xrpl/resource/Charge.h>
60#include <xrpl/resource/Consumer.h>
61#include <xrpl/resource/Disposition.h>
62#include <xrpl/resource/Fees.h>
63#include <xrpl/resource/Gossip.h>
64#include <xrpl/server/LoadFeeTrack.h>
65#include <xrpl/server/Manifest.h>
66#include <xrpl/server/NetworkOPs.h>
67#include <xrpl/shamap/SHAMap.h>
68#include <xrpl/shamap/SHAMapNodeID.h>
69#include <xrpl/tx/apply.h>
71#include <boost/algorithm/string/predicate.hpp>
72#include <boost/asio/bind_executor.hpp>
73#include <boost/asio/buffer.hpp>
74#include <boost/asio/completion_condition.hpp>
75#include <boost/asio/dispatch.hpp>
76#include <boost/asio/error.hpp>
77#include <boost/asio/strand.hpp>
78#include <boost/asio/write.hpp>
79#include <boost/beast/core/multi_buffer.hpp>
80#include <boost/beast/core/ostream.hpp>
81#include <boost/system/system_error.hpp>
83#include <google/protobuf/message.h>
107using namespace std::chrono_literals;
115constexpr std::chrono::milliseconds kPeerHighLatency{300};
120constexpr std::chrono::seconds kPeerTimerInterval{60};
178 <<
" vp reduce-relay base squelch enabled "
185 bool const inCluster{
cluster()};
210 if (
UInt256 ret; ret.parseHex(value))
222 if (
auto const iter = self->headers_.find(
"Closed-Ledger"); iter != self->headers_.end())
224 closed = parseLedgerHash(iter->value());
227 self->fail(
"Malformed handshake data (1)");
230 if (
auto const iter = self->headers_.find(
"Previous-Ledger"); iter != self->headers_.end())
232 previous = parseLedgerHash(iter->value());
235 self->fail(
"Malformed handshake data (2)");
238 if (previous && !closed)
239 self->fail(
"Malformed handshake data (3)");
244 self->closedLedgerHash_ = *closed;
246 self->previousLedgerHash_ = *previous;
255 self->doProtocolStart();
267 if (!self->socket_.is_open())
280 if (self->gracefulClose_)
282 if (self->detaching_)
284 if (!self->socket_.is_open())
287 auto validator = m->getValidatorKey();
288 if (validator && !self->squelch_.expireSquelch(*validator))
290 self->overlay_.reportOutboundTraffic(
292 static_cast<int>(m->getBuffer(self->compressionEnabled_).size()));
297 self->overlay_.reportOutboundTraffic(
299 static_cast<int>(m->getBuffer(self->compressionEnabled_).size()));
302 self->overlay_.reportOutboundTraffic(
304 static_cast<int>(m->getBuffer(self->compressionEnabled_).size()));
306 auto sendqSize = self->sendQueue_.size();
313 self->largeSendq_ = 0;
316 auto sink = self->journal_.debug();
320 sink << n <<
" sendq: " << sendqSize;
323 self->sendQueue_.push(m);
328 boost::asio::async_write(
330 boost::asio::buffer(self->sendQueue_.front()->getBuffer(self->compressionEnabled_)),
332 self->onWriteMessage(ec, bytesTransferred);
341 if (!self->txQueue_.empty())
343 protocol::TMHaveTransactions ht;
345 self->txQueue_, [&](
auto const& hash) { ht.add_hashes(hash.data(), hash.size()); });
346 JLOG(self->pJournal_.trace()) <<
"sendTxQueue " << self->txQueue_.size();
347 self->txQueue_.clear();
359 JLOG(self->pJournal_.warn()) <<
"addTxQueue exceeds the cap";
363 self->txQueue_.insert(hash);
364 JLOG(self->pJournal_.trace()) <<
"addTxQueue " << self->txQueue_.size();
372 auto removed = self->txQueue_.erase(hash);
373 JLOG(self->pJournal_.trace()) <<
"removeTxQueue " << removed;
382 self->usage_.disconnect(self->pJournal_))
390 bool expected =
false;
391 if (self->chargeDisconnectFired_.compare_exchange_strong(
392 expected,
true, std::memory_order_acq_rel))
394 self->overlay_.incPeerDisconnectCharges();
395 self->fail(
"charge: Resources");
406 auto const iter =
headers_.find(
"Crawl");
409 return boost::iequals(iter->value(),
"public");
435 ret[jss::inbound] =
true;
439 ret[jss::cluster] =
true;
451 if (
auto const nid =
headers_[
"Network-ID"]; !nid.empty())
454 ret[jss::load] =
usage_.balance();
473 if ((minSeq != 0) || (maxSeq != 0))
479 ret[jss::track] =
"diverged";
483 ret[jss::track] =
"unknown";
492 protocol::TMStatusChange lastStatus;
500 ret[jss::ledger] =
to_string(closedLedgerHash);
502 if (lastStatus.has_newstatus())
504 switch (lastStatus.newstatus())
506 case protocol::nsCONNECTING:
507 ret[jss::status] =
"connecting";
510 case protocol::nsCONNECTED:
511 ret[jss::status] =
"connected";
514 case protocol::nsMONITORING:
515 ret[jss::status] =
"monitoring";
518 case protocol::nsVALIDATING:
519 ret[jss::status] =
"validating";
522 case protocol::nsSHUTTING:
523 ret[jss::status] =
"shutting";
527 JLOG(
pJournal_.warn()) <<
"Unknown status: " << lastStatus.newstatus();
607 XRPL_ASSERT(
strand_.running_in_this_thread(),
"xrpl::PeerImp::close : strand in this thread");
628 JLOG(self->journal_.warn()) << n <<
" failed: " << reason;
637 XRPL_ASSERT(
strand_.running_in_this_thread(),
"xrpl::PeerImp::fail : strand in this thread");
650 strand_.running_in_this_thread(),
"xrpl::PeerImp::gracefulClose : strand in this thread");
651 XRPL_ASSERT(
socket_.is_open(),
"xrpl::PeerImp::gracefulClose : socket is open");
652 XRPL_ASSERT(!
gracefulClose_,
"xrpl::PeerImp::gracefulClose : socket is not closing");
657 stream_.async_shutdown(bind_executor(
666 timer_.expires_after(kPeerTimerInterval);
668 catch (boost::system::system_error
const& e)
670 JLOG(
journal_.error()) <<
"setTimer: " << e.code();
673 timer_.async_wait(bind_executor(
685 catch (boost::system::system_error
const&)
709 if (ec == boost::asio::error::operation_aborted)
713 JLOG(
journal_.error()) <<
"onTimer: " << ec.message();
720 fail(
"Large send queue");
726 ClockType::duration duration;
745 fail(
"Ping Timeout");
752 protocol::TMPing message;
753 message.set_type(protocol::TMPing::ptPING);
773 bool const shouldLog =
774 (ec != boost::asio::error::eof && ec != boost::asio::error::operation_aborted &&
775 !ec.message().contains(
"application data after close notify"));
779 JLOG(
journal_.debug()) <<
"onShutdown: " << ec.message();
790 XRPL_ASSERT(
readBuffer_.size() == 0,
"xrpl::PeerImp::doAccept : empty read buffer");
798 fail(
"makeSharedValue: Unexpected failure");
810 JLOG(
journal_.info()) <<
"Cluster name: " << *member;
822 !
overlay_.peerFinder().config().peerPrivate,
832 boost::asio::async_write(
835 boost::asio::transfer_all(),
844 if (ec == boost::asio::error::operation_aborted)
847 fail(
"onWriteResponse", ec);
851 if (writeBuffer->size() == bytesTransferred)
856 fail(
"Failed to write header");
886 app_.getValidators().forEachAvailable(
901 app_.getHashRouter(),
905 app_.getHashRouter().addSuppressionPeer(hash,
id_);
909 if (
auto m =
overlay_.getManifestsMessage())
924 if (ec == boost::asio::error::operation_aborted)
927 if (ec == boost::asio::error::eof)
934 fail(
"onReadMessage", ec);
940 stream <<
"onReadMessage: "
941 << (bytesTransferred > 0 ?
to_string(bytesTransferred) +
" bytes" :
"");
944 metrics_.recv.addMessage(bytesTransferred);
954 using namespace std::chrono_literals;
957 "invokeProtocolMessage",
963 fail(
"onReadMessage", ec);
973 if (bytesConsumed == 0)
984 self->onReadMessage(ec, bytesTransferred);
996 if (ec == boost::asio::error::operation_aborted)
999 fail(
"onWriteMessage", ec);
1002 if (
auto stream =
journal_.trace())
1004 stream <<
"onWriteMessage: "
1005 << (bytesTransferred > 0 ?
to_string(bytesTransferred) +
" bytes" :
"");
1008 metrics_.sent.addMessage(bytesTransferred);
1010 XRPL_ASSERT(!
sendQueue_.empty(),
"xrpl::PeerImp::onWriteMessage : non-empty send buffer");
1015 boost::asio::async_write(
1021 self->onWriteMessage(ec, bytesTransferred);
1028 stream_.async_shutdown(bind_executor(
1066 overlay_.reportInboundTraffic(category,
static_cast<int>(size));
1069 if ((type == MessageType::mtTRANSACTION || type == MessageType::mtHAVE_TRANSACTIONS ||
1070 type == MessageType::mtTRANSACTIONS ||
1083 JLOG(
journal_.trace()) <<
"onMessageBegin: " << type <<
" " << size <<
" " << uncompressedSize
1084 <<
" " << isCompressed;
1097 auto const s = m->list_size();
1119 if (m->type() == protocol::TMPing::ptPING)
1123 protocol::TMPing pong;
1124 pong.set_type(protocol::TMPing::ptPONG);
1126 pong.set_seq(m->seq());
1131 if (m->type() == protocol::TMPing::ptPONG && m->has_seq())
1170 for (
int i = 0; i < m->clusternodes().size(); ++i)
1172 protocol::TMClusterNode
const& node = m->clusternodes(i);
1175 if (node.has_nodename())
1176 name = node.nodename();
1186 app_.getCluster().update(*publicKey,
name, node.nodeload(), reportTime);
1190 int const loadSources = m->loadsources().size();
1191 if (loadSources != 0)
1194 gossip.
items.reserve(loadSources);
1195 for (
int i = 0; i < m->loadsources().size(); ++i)
1197 protocol::TMLoadSource
const& node = m->loadsources(i);
1202 gossip.
items.push_back(item);
1204 overlay_.resourceManager().importConsumers(
name(), gossip);
1208 auto const thresh =
app_.getTimeKeeper().now() - 90s;
1214 app_.getCluster().forEach([&fees, thresh](
ClusterNode const& status) {
1215 if (status.getReportTime() >= thresh)
1221 auto const index = fees.
size() / 2;
1223 clusterFee = fees[index];
1226 app_.getFeeTrack().setClusterFee(clusterFee);
1239 if (m->endpoints_v2().size() >= 1024)
1246 endpoints.
reserve(m->endpoints_v2().size());
1249 for (
auto const& tm : m->endpoints_v2())
1256 <<
"failed to parse incoming endpoint: {" << tm.endpoint() <<
"}";
1282 if (!endpoints.
empty())
1298 XRPL_ASSERT(eraseTxQueue != batch, (
"xrpl::PeerImp::handleTransaction : valid inputs"));
1302 if (
app_.getOPs().isNeedNetworkLedger())
1306 JLOG(
pJournal_.debug()) <<
"Ignoring incoming transaction: Need network ledger";
1315 UInt256 const txID = stx->getTransactionID();
1337 JLOG(
pJournal_.warn()) <<
"Ignoring Network relayed Tx containing "
1338 "tfInnerBatchTxn (handleTransaction).";
1347 if (!
app_.getHashRouter().shouldProcess(txID,
id_, flags, kTxInterval))
1353 JLOG(
pJournal_.debug()) <<
"Ignoring known bad tx " << txID;
1369 JLOG(
pJournal_.debug()) <<
"Got tx " << txID;
1371 bool checkSignature =
true;
1374 if (!m->has_deferred() || !m->deferred())
1383 if (!
app_.getValidationPublicKey())
1387 checkSignature =
false;
1391 if (
app_.getLedgerMaster().getValidatedLedgerAge() > 4min)
1393 JLOG(
pJournal_.trace()) <<
"No new transactions until synchronized";
1398 JLOG(
pJournal_.info()) <<
"Transaction queue is full";
1402 app_.getJobQueue().addJob(
1410 if (
auto peer = weak.lock())
1411 peer->checkTransaction(flags, checkSignature, stx, batch);
1421 JLOG(
pJournal_.warn()) <<
"Transaction invalid: " <<
strHex(m->rawtransaction())
1422 <<
". Exception: " << ex.
what();
1431 JLOG(
pJournal_.warn()) <<
"TMGetLedger: " << msg;
1433 auto const itype{m->itype()};
1436 if (itype < protocol::liBASE || itype > protocol::liTS_CANDIDATE)
1438 badData(
"Invalid ledger info type");
1445 return std::nullopt;
1448 if (itype == protocol::liTS_CANDIDATE)
1450 if (!m->has_ledgerhash())
1452 badData(
"Invalid TX candidate set, missing TX set hash");
1457 !m->has_ledgerhash() && !m->has_ledgerseq() && (!ltype || *ltype != protocol::ltCLOSED))
1459 badData(
"Invalid request");
1464 if (ltype && (*ltype < protocol::ltACCEPTED || *ltype > protocol::ltCLOSED))
1466 badData(
"Invalid ledger type");
1473 badData(
"Invalid ledger hash");
1478 if (m->has_ledgerseq())
1480 auto const ledgerSeq{m->ledgerseq()};
1483 using namespace std::chrono_literals;
1484 if (
app_.getLedgerMaster().getValidatedLedgerAge() <= 10s &&
1485 ledgerSeq >
app_.getLedgerMaster().getValidLedgerIndex() + 10)
1494 if (itype != protocol::liBASE)
1496 if (m->nodeids_size() <= 0)
1498 badData(
"Invalid ledger node IDs");
1505 "Requested number of ledger node IDs must be less than or equal to " +
1512 if (m->has_querytype() && m->querytype() != protocol::qtINDIRECT)
1514 badData(
"Invalid query type");
1519 if (m->has_querydepth())
1523 badData(
"Invalid query depth");
1530 app_.getJobQueue().addJob(
JtLedgerReq,
"RcvGetLedger", [weak, m, itype]() {
1531 auto peer = weak.
lock();
1536 bool tooManyNodeIds =
false;
1537 if (itype != protocol::liBASE)
1540 for (
auto const& nodeId : m->nodeids())
1547 tooManyNodeIds =
true;
1569 m->mutable_nodeids()->DeleteSubrange(
1570 static_cast<int>(nodeIDs.
size()),
1571 m->nodeids_size() -
static_cast<int>(nodeIDs.
size()));
1573 if (!m->has_requestcookie())
1578 peer->processLedgerRequest(m, std::move(nodeIDs));
1585 JLOG(
pJournal_.trace()) <<
"onMessage, TMProofPathRequest";
1595 if (
auto peer = weak.
lock())
1597 auto reply = peer->ledgerReplayMsgHandler_.processProofPathRequest(m);
1598 if (reply.has_error())
1600 if (reply.error() == protocol::TMReplyError::reBAD_REQUEST)
1642 JLOG(
pJournal_.trace()) <<
"onMessage, TMReplayDeltaRequest";
1652 if (
auto peer = weak.
lock())
1654 auto reply = peer->ledgerReplayMsgHandler_.processReplayDeltaRequest(m);
1655 if (reply.has_error())
1657 if (reply.error() == protocol::TMReplyError::reBAD_REQUEST)
1701 JLOG(
pJournal_.warn()) <<
"TMLedgerData: " << msg;
1707 badData(
"Invalid ledger hash");
1713 auto const ledgerSeq{m->ledgerseq()};
1714 if (m->type() == protocol::liTS_CANDIDATE)
1725 using namespace std::chrono_literals;
1726 if (
app_.getLedgerMaster().getValidatedLedgerAge() <= 10s &&
1727 ledgerSeq >
app_.getLedgerMaster().getValidLedgerIndex() + 10)
1736 if (m->type() < protocol::liBASE || m->type() > protocol::liTS_CANDIDATE)
1738 badData(
"Invalid ledger info type");
1743 if (m->has_error() &&
1744 (m->error() < protocol::reNO_LEDGER || m->error() > protocol::reBAD_REQUEST))
1746 badData(
"Invalid reply error");
1753 badData(
"Invalid Ledger/TXset nodes " +
std::to_string(m->nodes_size()));
1758 if (m->has_requestcookie())
1760 if (
auto peer =
overlay_.findPeerByShortID(m->requestcookie()))
1762 m->clear_requestcookie();
1769 auto const peerSupportsNodeDepth =
1772 MessageType messageType = MessageType::Unknown;
1773 for (
int i = 0; i < m->nodes_size(); ++i)
1775 auto* ledgerNode = m->mutable_nodes(i);
1779 if (ledgerNode->nodedata().empty())
1782 "Received node with empty data while relaying ledger data for " +
1788 MessageType msgType = MessageType::Unknown;
1789 if (m->type() == protocol::liBASE)
1791 if (ledgerNode->has_nodeid() || ledgerNode->has_id() || ledgerNode->has_depth())
1794 "Received liBASE message with node reference while relaying ledger "
1800 msgType = MessageType::Base;
1804 msgType = ledgerNode->has_nodeid() ? MessageType::Legacy : MessageType::Depth;
1806 if (messageType != MessageType::Unknown && messageType != msgType)
1809 "Received mixed mode message while relaying ledger data for " +
1814 messageType = msgType;
1816 if (peerSupportsNodeDepth || msgType != MessageType::Depth)
1820 !peerSupportsNodeDepth,
1821 "xrpl::PeerImp : relaying depth-format ledger data to pre-2.3 peer");
1822 switch (ledgerNode->reference_case())
1824 case protocol::TMLedgerNode::kId: {
1827 REACHABLE(
"xrpl::PeerImp : relay downgrade id to nodeid");
1828 ledgerNode->set_nodeid(ledgerNode->id());
1829 ledgerNode->clear_id();
1832 case protocol::TMLedgerNode::kDepth: {
1834 auto treeNode =
getTreeNode(ledgerNode->nodedata());
1838 "Unable to get tree node while relaying ledger data for " +
1848 "Unable to get node ID while relaying ledger data for " +
1854 REACHABLE(
"xrpl::PeerImp : relay downgrade depth to nodeid");
1855 ledgerNode->set_nodeid(nodeID->getRawString());
1856 ledgerNode->clear_depth();
1860 SOMETIMES(
true,
"xrpl::PeerImp : relay node has empty reference");
1862 "Empty node reference while relaying ledger data for " +
1874 JLOG(
pJournal_.info()) <<
"Unable to route TX/ledger data reply";
1882 if (m->type() == protocol::liTS_CANDIDATE)
1885 app_.getJobQueue().addJob(
JtTxnData,
"RcvPeerData", [weak, ledgerHash, m]() {
1886 if (
auto peer = weak.
lock())
1888 peer->app_.getInboundTransactions().gotData(ledgerHash, peer, m);
1901 protocol::TMProposeSet
const&
set = *m;
1910 JLOG(
pJournal_.warn()) <<
"Proposal: malformed";
1917 JLOG(
pJournal_.warn()) <<
"Proposal: malformed";
1926 auto const isTrusted =
app_.getValidators().trusted(publicKey);
1937 if (
app_.config().relayUntrustedProposals == -1)
1947 proposeHash, prevLedger,
set.proposeseq(), closeTime, publicKey.
slice(), sig);
1949 if (
auto [added, relayed] =
app_.getHashRouter().addSuppressionPeerWithStatus(suppression,
id_);
1955 overlay_.updateSlotAndSquelch(suppression, publicKey,
id_, protocol::mtPROPOSE_LEDGER);
1961 JLOG(
pJournal_.trace()) <<
"Proposal: duplicate";
1970 JLOG(
pJournal_.debug()) <<
"Proposal: Dropping untrusted (peer divergence)";
1974 if (!
cluster() &&
app_.getFeeTrack().isLoadedLocal())
1976 JLOG(
pJournal_.debug()) <<
"Proposal: Dropping untrusted (load)";
1981 JLOG(
pJournal_.trace()) <<
"Proposal: " << (isTrusted ?
"trusted" :
"untrusted");
1992 app_.getTimeKeeper().closeTime(),
1993 calcNodeID(
app_.getValidatorManifests().getMasterKey(publicKey))});
1996 app_.getJobQueue().addJob(
1998 if (
auto peer = weak.lock())
1999 peer->checkPropose(isTrusted, m,
proposal);
2006 JLOG(
pJournal_.trace()) <<
"Status: Change";
2008 if (!m->has_networktime())
2009 m->set_networktime(
app_.getTimeKeeper().now().time_since_epoch().count());
2013 if (!
lastStatus_.has_newstatus() || m->has_newstatus())
2020 protocol::NodeStatus
const status =
lastStatus_.newstatus();
2022 m->set_newstatus(status);
2026 if (m->newevent() == protocol::neLOST_SYNC)
2028 bool outOfSync{
false};
2042 JLOG(
pJournal_.debug()) <<
"Status: Out of sync";
2055 if (peerChangedLedgers)
2076 if (peerChangedLedgers)
2078 JLOG(
pJournal_.debug()) <<
"LCL is " << closedLedgerHash;
2082 JLOG(
pJournal_.debug()) <<
"Status: No ledger";
2086 if (m->has_firstseq() && m->has_lastseq())
2097 if (m->has_ledgerseq() &&
app_.getLedgerMaster().getValidatedLedgerAge() < 2min)
2105 if (m->has_newstatus())
2107 switch (m->newstatus())
2109 case protocol::nsCONNECTING:
2110 j[jss::status] =
"CONNECTING";
2112 case protocol::nsCONNECTED:
2113 j[jss::status] =
"CONNECTED";
2115 case protocol::nsMONITORING:
2116 j[jss::status] =
"MONITORING";
2118 case protocol::nsVALIDATING:
2119 j[jss::status] =
"VALIDATING";
2121 case protocol::nsSHUTTING:
2122 j[jss::status] =
"SHUTTING";
2127 if (m->has_newevent())
2129 switch (m->newevent())
2131 case protocol::neCLOSING_LEDGER:
2132 j[jss::action] =
"CLOSING_LEDGER";
2134 case protocol::neACCEPTED_LEDGER:
2135 j[jss::action] =
"ACCEPTED_LEDGER";
2137 case protocol::neSWITCHED_LEDGER:
2138 j[jss::action] =
"SWITCHED_LEDGER";
2140 case protocol::neLOST_SYNC:
2141 j[jss::action] =
"LOST_SYNC";
2146 if (m->has_ledgerseq())
2148 j[jss::ledger_index] = m->ledgerseq();
2151 if (m->has_ledgerhash())
2158 j[jss::ledger_hash] =
to_string(closedLedgerHash);
2161 if (m->has_networktime())
2166 if (m->has_firstseq() && m->has_lastseq())
2168 j[jss::ledger_index_min] =
json::UInt(m->firstseq());
2169 j[jss::ledger_index_max] =
json::UInt(m->lastseq());
2227 if (m->status() == protocol::tsHAVE)
2252 JLOG(
pJournal_.warn()) <<
"Ignored malformed " << messageType;
2258 auto const hash =
sha512Half(manifest, blobs, version);
2260 JLOG(
pJournal_.debug()) <<
"Received " << messageType;
2262 if (!
app_.getHashRouter().addSuppressionPeer(hash,
id_))
2264 JLOG(
pJournal_.debug()) << messageType <<
": received duplicate " << messageType;
2272 auto const applyResult =
app_.getValidators().applyListsAndBroadcast(
2279 app_.getHashRouter(),
2282 JLOG(
pJournal_.debug()) <<
"Processed " << messageType <<
" version " << version <<
" from "
2283 << (applyResult.publisherKey ?
strHex(*applyResult.publisherKey)
2284 :
"unknown or invalid publisher")
2285 <<
" with best result " <<
to_string(applyResult.bestDisposition());
2288 switch (applyResult.bestDisposition())
2299 applyResult.publisherKey,
2300 "xrpl::PeerImp::onValidatorListMessage : publisher key is "
2303 auto const& pubKey = *applyResult.publisherKey;
2309 iter->second < applyResult.sequence,
2310 "xrpl::PeerImp::onValidatorListMessage : lower sequence");
2323 applyResult.sequence && applyResult.publisherKey,
2324 "xrpl::PeerImp::onValidatorListMessage : nonzero sequence "
2325 "and set publisher key");
2328 "xrpl::PeerImp::onValidatorListMessage : maximum sequence");
2341 "xrpl::PeerImp::onValidatorListMessage : invalid best list "
2347 switch (applyResult.worstDisposition())
2384 "xrpl::PeerImp::onValidatorListMessage : invalid worst list "
2390 for (
auto const& [disp, count] : applyResult.dispositions)
2396 JLOG(
pJournal_.debug()) <<
"Applied " << count <<
" new " << messageType;
2400 JLOG(
pJournal_.debug()) <<
"Applied " << count <<
" expired " << messageType;
2404 JLOG(
pJournal_.debug()) <<
"Processed " << count <<
" future " << messageType;
2408 <<
"Ignored " << count <<
" " << messageType <<
"(s) with current sequence";
2412 <<
"Ignored " << count <<
" " << messageType <<
"(s) with future sequence";
2415 JLOG(
pJournal_.warn()) <<
"Ignored " << count <<
"stale " << messageType;
2418 JLOG(
pJournal_.warn()) <<
"Ignored " << count <<
" untrusted " << messageType;
2422 <<
"Ignored " << count <<
"unsupported version " << messageType;
2425 JLOG(
pJournal_.warn()) <<
"Ignored " << count <<
"invalid " << messageType;
2430 "xrpl::PeerImp::onValidatorListMessage : invalid list "
2442 if (m->version() < 2)
2445 <<
"ValidatorListCollection: received invalid validator list "
2456 JLOG(
pJournal_.warn()) <<
"ValidatorListCollection: Exception, " << e.
what();
2457 using namespace std::string_literals;
2465 if (m->validation().size() < 50)
2467 JLOG(
pJournal_.warn()) <<
"Validation: Too small";
2474 auto const closeTime =
app_.getTimeKeeper().closeTime();
2484 return calcNodeID(
app_.getValidatorManifests().getMasterKey(pk));
2487 .checkSignature =
false, .requireCanonicalOrder =
true});
2491 JLOG(
pJournal_.warn()) <<
"Validation: Exception, " << e.
what();
2495 val->setSeen(closeTime);
2499 app_.getValidations().parms(),
2500 app_.getTimeKeeper().closeTime(),
2502 val->getSeenTime()))
2504 JLOG(
pJournal_.trace()) <<
"Validation: Not current";
2512 auto const isTrusted =
app_.getValidators().trusted(val->getSignerPublic());
2523 if (
app_.config().relayUntrustedValidations == -1)
2529 auto [added, relayed] =
app_.getHashRouter().addSuppressionPeerWithStatus(key,
id_);
2539 key, val->getSignerPublic(),
id_, protocol::mtVALIDATION);
2546 JLOG(
pJournal_.trace()) <<
"Validation: duplicate";
2552 JLOG(
pJournal_.debug()) <<
"Dropping untrusted validation from diverged peer";
2554 else if (isTrusted || !
app_.getFeeTrack().isLoadedLocal())
2559 app_.getJobQueue().addJob(
2561 if (
auto peer = weak.
lock())
2562 peer->checkValidation(val, key, m);
2567 JLOG(
pJournal_.debug()) <<
"Dropping untrusted validation for load";
2572 JLOG(
pJournal_.warn()) <<
"Exception processing validation: " << e.
what();
2573 using namespace std::string_literals;
2581 protocol::TMGetObjectByHash
const& packet = *m;
2583 JLOG(
pJournal_.trace()) <<
"received TMGetObjectByHash " << packet.type() <<
" "
2584 << packet.objects_size();
2591 JLOG(
pJournal_.debug()) <<
"GetObject: Large send queue";
2595 if (packet.type() == protocol::TMGetObjectByHash::otFETCH_PACK)
2601 if (packet.type() == protocol::TMGetObjectByHash::otTRANSACTIONS)
2605 JLOG(
pJournal_.error()) <<
"TMGetObjectByHash: tx reduce-relay is disabled";
2612 if (
auto peer = weak.
lock())
2613 peer->doTransactions(m);
2618 if (packet.has_ledgerhash())
2622 JLOG(
pJournal_.debug()) <<
"GetObj: malformed ledgerhash from peer " <<
id_;
2633 <<
"GetObj: oversized request from peer " <<
id_ <<
" (" << packet.objects_size()
2643 bool const queued =
app_.getJobQueue().addJob(
JtLedgerReq,
"RcvGetObjByHash", [weak, m]() {
2644 auto peer = weak.
lock();
2649 peer->processGetObjectByHash(m);
2656 JLOG(peer->pJournal_.warn()) <<
"GetObj: handler threw: " << e.
what();
2664 JLOG(
pJournal_.warn()) <<
"GetObj: job queue refused request from peer " <<
id_;
2681 bool progress =
false;
2683 for (
int i = 0; i < packet.objects_size(); ++i)
2685 protocol::TMIndexedObject
const& obj = packet.objects(i);
2689 if (obj.has_ledgerseq())
2691 if (obj.ledgerseq() != pLSeq)
2693 if (pLDo && (pLSeq != 0))
2695 JLOG(
pJournal_.debug()) <<
"GetObj: Full fetch pack for " << pLSeq;
2697 pLSeq = obj.ledgerseq();
2698 pLDo = !
app_.getLedgerMaster().haveLedger(pLSeq);
2702 JLOG(
pJournal_.debug()) <<
"GetObj: Late fetch pack for " << pLSeq;
2715 app_.getLedgerMaster().addFetchPack(
2721 if (pLDo && (pLSeq != 0))
2723 JLOG(
pJournal_.debug()) <<
"GetObj: Partial fetch pack for " << pLSeq;
2725 if (packet.type() == protocol::TMGetObjectByHash::otFETCH_PACK)
2726 app_.getLedgerMaster().gotFetchPack(progress, pLSeq);
2733 protocol::TMGetObjectByHash
const& packet = *m;
2735 protocol::TMGetObjectByHash reply;
2736 reply.set_query(
false);
2737 reply.set_type(packet.type());
2739 if (packet.has_ledgerhash())
2741 reply.set_ledgerhash(packet.ledgerhash());
2752 int const requested = packet.objects_size();
2755 for (
int i = 0; i < iterLimit; ++i)
2757 auto const& obj = packet.objects(i);
2764 std::uint32_t const seq{obj.has_ledgerseq() ? obj.ledgerseq() : 0};
2765 auto const nodeObject =
app_.getNodeStore().fetchNodeObject(hash, seq);
2769 protocol::TMIndexedObject& newObj = *reply.add_objects();
2770 newObj.set_hash(hash.
begin(), hash.
size());
2771 auto const& data = nodeObject->getData();
2772 newObj.set_data(data.data(), data.size());
2773 if (obj.has_nodeid())
2774 newObj.set_index(obj.nodeid());
2775 if (obj.has_ledgerseq())
2776 newObj.set_ledgerseq(obj.ledgerseq());
2787 "processed get object by hash request");
2789 JLOG(
pJournal_.trace()) <<
"GetObj: " << reply.objects_size() <<
" of " << requested;
2798 JLOG(
pJournal_.error()) <<
"TMHaveTransactions: tx reduce-relay is disabled";
2805 if (
auto peer = weak.
lock())
2806 peer->handleHaveTransactions(m);
2813 protocol::TMGetObjectByHash tmBH;
2814 tmBH.set_type(protocol::TMGetObjectByHash_ObjectType_otTRANSACTIONS);
2815 tmBH.set_query(
true);
2817 JLOG(
pJournal_.trace()) <<
"received TMHaveTransactions " << m->hashes_size();
2823 JLOG(
pJournal_.error()) <<
"TMHaveTransactions with invalid hash size";
2830 auto txn =
app_.getMasterTransaction().fetchFromCache(hash);
2832 JLOG(
pJournal_.trace()) <<
"checking transaction " << (bool)txn;
2836 JLOG(
pJournal_.debug()) <<
"adding transaction to request";
2838 auto obj = tmBH.add_objects();
2839 obj->set_hash(hash.
data(), hash.
size());
2850 JLOG(
pJournal_.trace()) <<
"transaction request object is " << tmBH.objects_size();
2852 if (tmBH.objects_size() > 0)
2861 JLOG(
pJournal_.error()) <<
"TMTransactions: tx reduce-relay is disabled";
2868 JLOG(
pJournal_.error()) <<
"TMTransactions: transaction list too large";
2873 JLOG(
pJournal_.trace()) <<
"received TMTransactions " << m->transactions_size();
2875 overlay_.addTxMetrics(m->transactions_size());
2881 m->mutable_transactions(i), [](protocol::TMTransaction*) {}),
2891 if (!m->has_validatorpubkey())
2896 auto validator = m->validatorpubkey();
2906 if (key == self->app_.getValidationPublicKey())
2908 JLOG(self->pJournal_.debug())
2909 <<
"onMessage: TMSquelch discarding validator's squelch " << slice;
2913 std::uint32_t const duration = m->has_squelchduration() ? m->squelchduration() : 0;
2916 self->squelch_.removeSquelch(key);
2923 JLOG(self->pJournal_.debug())
2924 <<
"onMessage: TMSquelch " << slice <<
" " << self->id() <<
" " << duration;
2935 (void)lockedRecentLock;
2949 if (
app_.getFeeTrack().isLoadedLocal() ||
2950 (
app_.getLedgerMaster().getValidatedLedgerAge() > 40s) ||
2951 (
app_.getJobQueue().getJobCount(
JtPack) > 10))
2953 JLOG(
pJournal_.info()) <<
"Too busy to make fetch pack";
2959 JLOG(
pJournal_.warn()) <<
"FetchPack hash size malformed";
2970 auto const pap = &
app_;
2971 app_.getJobQueue().addJob(
JtPack,
"MakeFetchPack", [pap, weak, packet, hash, elapsed]() {
2972 pap->getLedgerMaster().makeFetchPack(weak, packet, hash, elapsed);
2979 protocol::TMTransactions reply;
2981 JLOG(
pJournal_.trace()) <<
"received TMGetObjectByHash requesting tx "
2982 << packet->objects_size();
2986 JLOG(
pJournal_.error()) <<
"doTransactions, invalid number of hashes";
2993 auto const& obj = packet->objects(i);
3003 auto txn =
app_.getMasterTransaction().fetchFromCache(hash);
3008 <<
"doTransactions, transaction not found " <<
Slice(hash.
data(), hash.
size());
3014 auto tx = reply.add_transactions();
3015 auto sttx = txn->getSTransaction();
3017 tx->set_rawtransaction(s.
data(), s.
size());
3020 tx->set_receivetimestamp(
app_.getTimeKeeper().now().time_since_epoch().count());
3021 tx->set_deferred(txn->getSubmitResult().queued);
3024 if (reply.transactions_size() > 0)
3031 bool checkSignature,
3058 JLOG(
pJournal_.warn()) <<
"Ignoring Network relayed Tx containing "
3059 "tfInnerBatchTxn (checkSignature).";
3066 if (stx->isFieldPresent(sfLastLedgerSequence) &&
3067 (stx->getFieldU32(sfLastLedgerSequence) <
app_.getLedgerMaster().getValidLedgerIndex()))
3069 JLOG(
pJournal_.info()) <<
"Marking transaction " << stx->getTransactionID()
3070 <<
"as BAD because it's expired";
3084 "xrpl::PeerImp::checkTransaction Transaction created "
3088 JLOG(
pJournal_.debug()) <<
"Processing " << (batch ?
"batch" :
"unsolicited")
3089 <<
" pseudo-transaction tx " << tx->getID();
3091 app_.getMasterTransaction().canonicalize(&tx);
3093 auto const toSkip =
app_.getHashRouter().shouldRelay(tx->getID());
3097 <<
"Passing skipped pseudo pseudo-transaction tx " << tx->getID();
3098 app_.getOverlay().relay(tx->getID(), {}, *toSkip);
3102 JLOG(
pJournal_.debug()) <<
"Charging for pseudo-transaction tx " << tx->getID();
3113 auto const& validatedRules =
app_.getLedgerMaster().getValidatedRules();
3114 if (
auto [
valid, validReason] =
3118 if (!validReason.empty())
3120 JLOG(
pJournal_.debug()) <<
"Exception checking transaction: " << validReason;
3131 if (validatedRules.enabled(fixCleanup3_4_0) ||
3132 (!stx->isFieldPresent(sfSponsorSignature) &&
3133 !stx->isFieldPresent(sfCounterpartySignature)))
3151 if (!reason.
empty())
3153 JLOG(
pJournal_.debug()) <<
"Exception checking transaction: " << reason;
3165 JLOG(
pJournal_.warn()) <<
"Exception in " << __func__ <<
": " << ex.
what();
3167 using namespace std::string_literals;
3179 JLOG(
pJournal_.trace()) <<
"Checking " << (isTrusted ?
"trusted" :
"UNTRUSTED") <<
" proposal";
3181 XRPL_ASSERT(packet,
"xrpl::PeerImp::checkPropose : non-null packet");
3185 std::string const desc{
"Proposal fails sig check"};
3195 relay =
app_.getOPs().processTrustedProposal(peerPos);
3199 relay =
app_.config().relayUntrustedProposals == 1 ||
cluster();
3210 if (!haveMessage.empty())
3215 std::move(haveMessage),
3216 protocol::mtPROPOSE_LEDGER);
3227 if (!val->isValid())
3229 std::string const desc{
"Validation forwarded by peer is invalid"};
3244 auto haveMessage =
overlay_.relay(*packet, key, val->getSignerPublic());
3245 if (!haveMessage.empty())
3248 key, val->getSignerPublic(), std::move(haveMessage), protocol::mtVALIDATION);
3254 JLOG(
pJournal_.trace()) <<
"Exception processing validation: " << ex.
what();
3255 using namespace std::string_literals;
3270 if (p->hasTxSet(rootHash) && p.get() != skip)
3272 auto score = p->getScore(
true);
3273 if (!ret || (score > retScore))
3298 if (p->hasLedger(ledgerHash, ledger) && p.get() != skip)
3300 auto score = p->getScore(
true);
3301 if (!ret || (score > retScore))
3315 protocol::TMLedgerData& ledgerData)
3317 JLOG(
pJournal_.trace()) <<
"sendLedgerBase: Base data";
3320 addRaw(ledger->header(), s);
3323 auto const& stateMap{ledger->stateMap()};
3329 stateMap.serializeRoot(
root);
3330 ledgerData.add_nodes()->set_nodedata(
root.getDataPtr(),
root.getLength());
3334 auto const& txMap{ledger->txMap()};
3339 txMap.serializeRoot(
root);
3340 ledgerData.add_nodes()->set_nodedata(
root.getDataPtr(),
root.getLength());
3352 JLOG(
pJournal_.trace()) <<
"getLedger: Ledger";
3356 if (m->has_ledgerhash())
3360 ledger =
app_.getLedgerMaster().getLedgerByHash(ledgerHash);
3363 JLOG(
pJournal_.trace()) <<
"getLedger: Don't have ledger with hash " << ledgerHash;
3365 if (m->has_querytype() && !m->has_requestcookie())
3369 overlay_, ledgerHash, m->has_ledgerseq() ? m->ledgerseq() : 0,
this))
3371 m->set_requestcookie(
id());
3373 JLOG(
pJournal_.debug()) <<
"getLedger: Request relayed to peer";
3377 JLOG(
pJournal_.trace()) <<
"getLedger: Failed to find peer to relay request";
3381 else if (m->has_ledgerseq())
3384 if (m->ledgerseq() <
app_.getLedgerMaster().getEarliestFetch())
3386 JLOG(
pJournal_.debug()) <<
"getLedger: Early ledger sequence request";
3390 ledger =
app_.getLedgerMaster().getLedgerBySeq(m->ledgerseq());
3394 <<
"getLedger: Don't have ledger with sequence " << m->ledgerseq();
3398 else if (m->has_ltype() && m->ltype() == protocol::ltCLOSED)
3400 ledger =
app_.getLedgerMaster().getClosedLedger();
3406 auto const ledgerSeq{ledger->header().seq};
3407 if (m->has_ledgerseq())
3409 if (ledgerSeq != m->ledgerseq())
3412 if (!m->has_requestcookie())
3416 JLOG(
pJournal_.warn()) <<
"getLedger: Invalid ledger sequence " << ledgerSeq;
3419 else if (ledgerSeq <
app_.getLedgerMaster().getEarliestFetch())
3422 JLOG(
pJournal_.debug()) <<
"getLedger: Early ledger sequence request " << ledgerSeq;
3427 JLOG(
pJournal_.debug()) <<
"getLedger: Unable to find ledger";
3436 JLOG(
pJournal_.trace()) <<
"getTxSet: TX set";
3442 if (m->has_querytype() && !m->has_requestcookie())
3447 m->set_requestcookie(
id());
3449 JLOG(
pJournal_.debug()) <<
"getTxSet: Request relayed";
3453 JLOG(
pJournal_.debug()) <<
"getTxSet: Failed to find relay peer";
3458 JLOG(
pJournal_.debug()) <<
"getTxSet: Failed to find TX set";
3472 SHAMap const* map{
nullptr};
3473 protocol::TMLedgerData ledgerData;
3474 bool fatLeaves{
true};
3475 auto const itype{m->itype()};
3477 if (itype == protocol::liTS_CANDIDATE)
3479 if (sharedMap =
getTxSet(m); !sharedMap)
3481 map = sharedMap.
get();
3484 ledgerData.set_ledgerseq(0);
3485 ledgerData.set_ledgerhash(m->ledgerhash());
3486 ledgerData.set_type(protocol::liTS_CANDIDATE);
3487 if (m->has_requestcookie())
3488 ledgerData.set_requestcookie(m->requestcookie());
3497 JLOG(
pJournal_.debug()) <<
"processLedgerRequest: Large send queue";
3500 if (
app_.getFeeTrack().isLoadedLocal() && !
cluster())
3502 JLOG(
pJournal_.debug()) <<
"processLedgerRequest: Too busy";
3510 auto const ledgerHash{ledger->header().hash};
3511 ledgerData.set_ledgerhash(ledgerHash.begin(), ledgerHash.size());
3512 ledgerData.set_ledgerseq(ledger->header().seq);
3513 ledgerData.set_type(itype);
3514 if (m->has_requestcookie())
3515 ledgerData.set_requestcookie(m->requestcookie());
3519 case protocol::liBASE:
3523 case protocol::liTX_NODE:
3524 map = &ledger->txMap();
3529 case protocol::liAS_NODE:
3530 map = &ledger->stateMap();
3532 <<
"processLedgerRequest: Account state map hash " <<
to_string(map->
getHash());
3537 JLOG(
pJournal_.error()) <<
"processLedgerRequest: Invalid ledger info type";
3544 JLOG(
pJournal_.warn()) <<
"processLedgerRequest: Unable to find map";
3549 if (!nodeIDs.
empty())
3552 auto const queryDepth{m->has_querydepth() ? m->querydepth() : defaultDepth};
3558 for (
auto const& nodeID : nodeIDs)
3567 if (map->
getNodeFat(nodeID, data, fatLeaves, queryDepth))
3570 <<
"processLedgerRequest: getNodeFat got " << data.size() <<
" nodes";
3572 for (
auto const& d : data)
3577 protocol::TMLedgerNode* node{ledgerData.add_nodes()};
3578 node->set_nodedata(d.data.data(), d.data.size());
3583 if (!useLedgerNodeDepth)
3585 node->set_nodeid(d.nodeID.getRawString());
3589 REACHABLE(
"xrpl::PeerImp : emit leaf depth in reply");
3590 node->set_depth(d.nodeID.getDepth());
3594 REACHABLE(
"xrpl::PeerImp : emit inner id in reply");
3595 node->set_id(d.nodeID.getRawString());
3601 JLOG(
pJournal_.warn()) <<
"processLedgerRequest: getNodeFat returns false";
3609 case protocol::liBASE:
3611 info =
"Ledger base";
3614 case protocol::liTX_NODE:
3618 case protocol::liAS_NODE:
3622 case protocol::liTS_CANDIDATE:
3623 info =
"TS candidate";
3631 if (!m->has_ledgerhash())
3632 info +=
", no hash specified";
3635 <<
"processLedgerRequest: getNodeFat with nodeId " << nodeID
3636 <<
" and ledger info type " << info <<
" throws exception: " << e.
what();
3640 JLOG(
pJournal_.info()) <<
"processLedgerRequest: Got request for " << m->nodeids_size()
3641 <<
" node IDs at depth " << queryDepth <<
", return "
3642 << ledgerData.nodes_size() <<
" nodes";
3645 if (ledgerData.nodes_size() == 0)
3678 int const missed =
std::max(0, requested - found);
3679 int const billableMisses =
std::min(missed, billable);
3680 int const billableHits = billable - billableMisses;
3703 static int const kSpRandomMax = 9999;
3707 static int const kSpHaveItem = 10000;
3712 static int const kSpLatency = 30;
3715 static int const kSpNoLatency = 8000;
3717 int score =
randInt(kSpRandomMax);
3720 score += kSpHaveItem;
3730 score -= latency->count() * kSpLatency;
3734 score -= kSpNoLatency;
3744 return latency_ >= kPeerHighLatency;
3750 using namespace std::chrono_literals;
3758 if (timeElapsedInSecs >= 1s)
3760 auto const avgBytes =
accumBytes_ / timeElapsedInSecs.count();
A version-independent IP address and port combination.
static std::optional< Endpoint > fromStringChecked(std::string const &s)
Create an Endpoint from a string.
static Endpoint fromString(std::string const &s)
static BaseUInt fromRaw(Container const &c)
static constexpr std::size_t size()
static std::size_t messageSize(::google::protobuf::Message const &message)
std::chrono::time_point< NetClock > time_point
std::chrono::duration< rep, period > duration
Child(OverlayImpl &overlay)
void forEach(UnaryFunc &&f) const
std::uint64_t rollingAvgBytes_
boost::circular_buffer< std::uint64_t > rollingAvg_
void addMessage(std::uint64_t bytes)
std::uint64_t averageBytes() const
ClockType::time_point intervalStart_
std::uint64_t totalBytes() const
std::uint64_t accumBytes_
std::uint64_t totalBytes_
void checkTracking(std::uint32_t validationSeq)
Check if the peer is tracking.
void onTimer(boost::system::error_code const &ec)
std::optional< std::chrono::milliseconds > latency_
std::string getVersion() const
Return the version of xrpld that the peer is running, if reported.
void onMessage(std::shared_ptr< protocol::TMManifests > const &m)
ProtocolVersion protocol_
bool txReduceRelayEnabled_
void checkTransaction(HashRouterFlags flags, bool checkSignature, std::shared_ptr< STTx const > const &stx, bool batch)
void handleTransaction(std::shared_ptr< protocol::TMTransaction > const &m, bool eraseTxQueue, bool batch)
Called from onMessage(TMTransaction(s)).
bool hasTxSet(UInt256 const &hash) const override
void addLedger(UInt256 const &hash, std::scoped_lock< std::mutex > const &lockedRecentLock)
bool txReduceRelayEnabled() const override
compression::Compressed Compressed
boost::beast::http::fields const & headers_
Compressed compressionEnabled_
ClockType::duration uptime() const
ClockType::time_point const creationTime_
boost::circular_buffer< UInt256 > recentTxSets_
std::shared_ptr< peer_finder::Slot > const slot_
void onWriteMessage(ErrorCode ec, std::size_t bytesTransferred)
boost::beast::multi_buffer readBuffer_
void onShutdown(ErrorCode ec)
boost::circular_buffer< UInt256 > recentLedgers_
std::string const & fingerprint() const override
void checkValidation(std::shared_ptr< STValidation > const &val, UInt256 const &key, std::shared_ptr< protocol::TMValidation > const &packet)
bool ledgerReplayEnabled_
void sendTxQueue() override
Send aggregated transactions' hashes.
struct xrpl::PeerImp::@337373043150231020277011015352151251117171316327 metrics_
beast::ip::Endpoint const remoteAddress_
reduce_relay::Squelch< UptimeClock > squelch_
PeerImp(PeerImp const &)=delete
void cycleStatus() override
std::shared_mutex nameMutex_
std::string domain() const
std::atomic< Tracking > tracking_
void ledgerRange(std::uint32_t &minSeq, std::uint32_t &maxSeq) const override
HashMap< PublicKey, std::size_t > publisherListSequences_
void addTxQueue(UInt256 const &hash) override
Add transaction's hash to the transactions' hashes queue.
void onReadMessage(ErrorCode ec, std::size_t bytesTransferred)
beast::WrappedSink pSink_
PublicKey const publicKey_
void processGetObjectByHash(std::shared_ptr< protocol::TMGetObjectByHash > const &m)
Process a generic-query TMGetObjectByHash message.
bool cluster() const override
Returns true if this connection is a member of the cluster.
static resource::Charge computeGetObjectByHashFee(int const requested, int const found)
Compute the per-message resource charge for a TMGetObjectByHash request based on how much work was ac...
std::queue< std::shared_ptr< Message > > sendQueue_
void checkPropose(bool isTrusted, std::shared_ptr< protocol::TMProposeSet > const &packet, RCLCxPeerPos peerPos)
std::chrono::steady_clock ClockType
std::unique_ptr< StreamType > streamPtr_
resource::Consumer usage_
void sendLedgerBase(std::shared_ptr< Ledger const > const &ledger, protocol::TMLedgerData &ledgerData)
bool hasLedger(UInt256 const &hash, std::uint32_t seq) const override
beast::Journal const pJournal_
ClockType::time_point trackingTime_
static std::string makePrefix(std::string const &fingerprint)
void handleHaveTransactions(std::shared_ptr< protocol::TMHaveTransactions > const &m)
Handle protocol message with hashes of transactions that have not been relayed by an upstream node do...
void onMessageUnknown(std::uint16_t type)
void processLedgerRequest(std::shared_ptr< protocol::TMGetLedger > const &m, std::vector< SHAMapNodeID > nodeIDs)
Tracking
Whether the peer's view of the ledger converges or diverges from ours.
protocol::TMStatusChange lastStatus_
LedgerReplayMsgHandler ledgerReplayMsgHandler_
boost::asio::basic_waitable_timer< std::chrono::steady_clock > WaitableTimer
void fail(std::string const &reason)
std::unique_ptr< LoadEvent > loadEvent_
void send(std::shared_ptr< Message > const &m) override
void onMessageBegin(std::uint16_t type, std::shared_ptr<::google::protobuf::Message > const &m, std::size_t size, std::size_t uncompressedSize, bool isCompressed)
UInt256 closedLedgerHash_
void doFetchPack(std::shared_ptr< protocol::TMGetObjectByHash > const &packet)
bool supportsFeature(ProtocolFeature f) const override
std::shared_ptr< Ledger const > getLedger(std::shared_ptr< protocol::TMGetLedger > const &m)
boost::asio::strand< boost::asio::executor > strand_
int getScore(bool haveItem) const override
beast::Journal const journal_
bool isHighLatency() const override
void cancelTimer() noexcept
bool hasRange(std::uint32_t uMin, std::uint32_t uMax) override
std::shared_ptr< SHAMap const > getTxSet(std::shared_ptr< protocol::TMGetLedger > const &m) const
std::optional< std::uint32_t > lastPingSeq_
boost::system::error_code ErrorCode
bool crawl() const
Returns true if this connection will publicly share its IP address.
void onValidatorListMessage(std::string const &messageType, std::string const &manifest, std::uint32_t version, std::vector< ValidatorBlobInfo > const &blobs)
Peer::ID id() const override
void charge(resource::Charge const &fee, std::string const &context) override
Adjust this peer's load balance based on the type of load imposed.
void removeTxQueue(UInt256 const &hash) override
Remove transaction's hash from the transactions' hashes queue.
void doTransactions(std::shared_ptr< protocol::TMGetObjectByHash > const &packet)
Process peer's request to send missing transactions.
std::shared_ptr< peer_finder::Slot > const & slot()
UInt256 previousLedgerHash_
ClockType::time_point lastPingTime_
json::Value json() override
void onMessageEnd(std::uint16_t type, std::shared_ptr<::google::protobuf::Message > const &m)
std::uint32_t ID
Uniquely identifies a peer.
Slice slice() const noexcept
A peer's signed, proposed position for use in RCLConsensus.
PublicKey const & publicKey() const
Public key of peer that sent the proposal.
ConsensusProposal< NodeID, UInt256, UInt256 > Proposal
UInt256 const & suppressionID() const
Unique id used by hash router to suppress duplicates.
bool checkSign() const
Verify the signing hash of the proposal.
bool getNodeFat(SHAMapNodeID const &wanted, std::vector< SHAMapNodeData > &data, bool fatLeaves, std::uint32_t depth) const
SHAMapHash getHash() const
void const * getDataPtr() const
std::size_t size() const noexcept
void const * data() const noexcept
An immutable linear range of bytes.
static Category attribute(Category cat, IsFromCluster isFromCluster)
Limits a category to what the sender can report.
static Category categorize(::google::protobuf::Message const &message, protocol::MessageType type, bool inbound)
Given a protocol message, determine which traffic category it belongs to.
static void sendValidatorList(Peer &peer, std::uint64_t peerSequence, PublicKey const &publisherKey, std::size_t maxSequence, std::uint32_t rawVersion, std::string const &rawManifest, std::map< std::size_t, ValidatorBlobInfo > const &blobInfos, HashRouter &hashRouter, beast::Journal j)
static std::vector< ValidatorBlobInfo > parseBlobs(std::uint32_t version, json::Value const &body)
Pull the blob/signature/manifest information out of the appropriate Json body fields depending on the...
An endpoint that consumes resources.
T duration_cast(T... args)
T emplace_back(T... args)
@ Object
object value (collection of name/value pairs).
TER valid(STTx const &tx, ReadView const &view, AccountID const &src, beast::Journal j)
auto measureDurationAndLog(Func &&func, std::string_view actionDescription, std::chrono::duration< Rep, Period > maxDelay, beast::Journal const &journal)
static constexpr std::size_t kMaxTxQueueSize
static constexpr auto kIdled
Charge const kFeeMalformedData
Charge const kFeeHeavyBurdenPeer
Charge const kFeeUselessData
Charge const kFeeRequestNoReply
Charge const kFeeMalformedRequest
Schedule of fees charged for imposing load on the server.
Charge const kFeeInvalidSignature
Charge const kFeeInvalidData
Charge const kFeeModerateBurdenPeer
Charge const kFeeTrivialPeer
static constexpr auto kSoftMaxReplyNodes
The soft cap on the number of ledger entries in a single reply.
static constexpr auto kBandMediumMax
static constexpr auto kTargetSendQueue
How many messages we consider reasonable sustained on a send queue.
static constexpr auto kCostPerLookupHit
Cost of one cache-hit lookup.
static constexpr auto kCostBandLarge
static constexpr auto kFreeObjectsPerRequest
TMGetObjectByHash differential pricing.
static constexpr auto kHardMaxReplyNodes
The hard cap on the number of ledger entries in a single reply.
static constexpr std::uint32_t kConvergedLedgerLimit
How many ledgers off a server can be and we will still consider it converged.
static constexpr auto kDropSendQueue
How many messages on a send queue before we refuse queries.
static constexpr auto kMaxQueryDepth
The maximum number of levels to search.
static constexpr auto kSendqIntervals
How many timer intervals a sendq has to stay large before we disconnect.
static constexpr auto kBandSmallMax
Cutoffs that decide which size band a request falls into.
constexpr std::size_t kReadBufferBytes
Size of buffer used to read from the socket.
static constexpr auto kCostBandMedium
static constexpr auto kCostPerLookupMiss
Cost of one node-store miss, in units of kCostPerLookupHit.
static constexpr std::uint32_t kDivergedLedgerLimit
How many ledgers off a server has to be before we consider it diverged.
static constexpr auto kCostBandSmall
Size-band surcharges.
static constexpr auto kSendQueueLogFreq
How often to log send queue size.
Use hash_* containers for keys that do not need a cryptographically secure hashing algorithm.
bool set(T &target, std::string const &name, Section const §ion)
Set a value from a configuration Section If the named value is not found or doesn't parse as a T,...
std::string base64Decode(std::string_view data)
constexpr FlagValue tfInnerBatchTxn
bool isCurrent(ValidationParms const &p, NetClock::time_point now, NetClock::time_point signTime, NetClock::time_point seenTime)
Whether a validation is still current.
@ Valid
Signature and local checks are good / passed.
std::optional< AccountID > parseBase58(std::string const &s)
Parse AccountID from checked, base58 string.
std::uint32_t LedgerIndex
A ledger index.
Stopwatch & stopwatch()
Returns an instance of a wall clock.
std::string strHex(FwdIt begin, FwdIt end)
std::pair< Validity, std::string > checkValidity(HashRouter &router, STTx const &tx, Rules const &rules)
Checks transaction signature and local checks.
HttpResponseType makeResponse(bool crawlPublic, HttpRequestType const &req, beast::ip::Address publicIp, beast::ip::Address remoteIp, UInt256 const &sharedValue, std::optional< std::uint32_t > networkID, ProtocolVersion protocol, Application &app)
Make http response.
UInt256 proposalUniqueId(UInt256 const &proposeHash, UInt256 const &previousLedger, std::uint32_t proposeSeq, NetClock::time_point closeTime, Slice const &publicKey, Slice const &signature)
Calculate a unique identifier for a signed proposal.
@ UnsupportedVersion
List version is not supported.
@ Expired
List is expired, but has the largest non-pending sequence seen so far.
@ SameSequence
Same sequence as current list.
@ KnownSequence
Future sequence already seen.
@ Pending
List will be valid in the future.
@ Invalid
Invalid format or signature.
@ Untrusted
List signed by untrusted publisher key.
@ Stale
Trusted publisher key, but seq is too old.
std::string toBase58(AccountID const &v)
Convert AccountID to base58 checked string.
static constexpr char kFeatureLedgerReplay[]
Number root(Number f, unsigned d)
static std::shared_ptr< PeerImp > getPeerWithTree(OverlayImpl &ov, UInt256 const &rootHash, PeerImp const *skip)
constexpr Dest safeCast(Src s) noexcept
std::optional< SHAMapNodeID > getSHAMapNodeID(protocol::TMLedgerNode const &ledgerNode, SHAMapTreeNode const &treeNode)
Extracts or reconstructs the SHAMapNodeID from a ledger node proto message.
std::string to_string(BaseUInt< Bits, Tag > const &a)
std::optional< KeyType > publicKeyType(Slice const &slice)
Returns the type of public key.
SHAMapTreeNodePtr getTreeNode(std::string_view data)
Deserializes a SHAMapTreeNode from wire format data.
static constexpr char kFeatureTxrr[]
Slice makeSlice(std::array< T, N > const &a)
std::pair< std::size_t, boost::system::error_code > invokeProtocolMessage(Buffers const &buffers, Handler &handler, std::size_t &hint)
Calls the handler for up to one protocol message in the passed buffers.
void forceValidity(HashRouter &router, UInt256 const &txid, Validity validity)
Sets the validity of a given transaction in the cache.
void addRaw(LedgerHeader const &, Serializer &, bool includeHash=false)
boost::beast::http::request< boost::beast::http::dynamic_body > HttpRequestType
NodeID calcNodeID(PublicKey const &)
Calculate the 160-bit node ID from a node public key.
std::optional< UInt256 > makeSharedValue(StreamType &ssl, beast::Journal journal)
Computes a shared value based on the SSL connection state.
constexpr ProtocolVersion makeProtocol(std::uint16_t major, std::uint16_t minor)
std::optional< SHAMapNodeID > deserializeSHAMapNodeID(void const *data, std::size_t size)
Return an object representing a serialized SHAMap Node ID.
std::string getFingerprint(beast::ip::Endpoint const &address, std::optional< PublicKey > const &publicKey=std::nullopt, std::optional< std::string > const &id=std::nullopt)
bool peerFeatureEnabled(Headers const &request, std::string const &feature, std::string value, bool config)
Check if a feature should be enabled for a peer.
@ Malformed
Protocol-level violation; no honest peer would produce this.
@ BadData
Peer reported has_error() (legitimate "cannot fulfill" signal).
std::pair< std::uint16_t, std::uint16_t > ProtocolVersion
Represents a particular version of the peer-to-peer protocol.
static bool stringIsUInt256Sized(std::string const &pBuffStr)
static constexpr char kFeatureCompr[]
static std::shared_ptr< PeerImp > getPeerWithLedger(OverlayImpl &ov, UInt256 const &ledgerHash, LedgerIndex ledger, PeerImp const *skip)
std::string protocolMessageName(int type)
Returns the name of a protocol message given its type.
bool isPseudoTx(STObject const &tx)
Check whether a transaction is a pseudo-transaction.
static constexpr char kFeatureVprr[]
Sha512HalfHasher::result_type sha512Half(Args const &... args)
Returns the SHA512-Half of a series of objects.
constexpr bool any(HashRouterFlags flags)
T shared_from_this(T... args)
Options controlling deserialization of a STValidation.
Describes a single consumer.
beast::ip::Endpoint address
Data format for exchanging consumption information across peers.
std::vector< Item > items