From 4ab295041004a719b239b1b76ca00fbb41b305de Mon Sep 17 00:00:00 2001 From: Ed Hennis Date: Mon, 13 Apr 2026 18:40:46 -0400 Subject: [PATCH] Reduce duplicate peer traffic for ledger data (#5126) Improve job queue collision checks and logging - Improve logging related to ledger acquisition and operating mode changes - Class "CanProcess" to keep track of processing of distinct items - Drop duplicate outgoing TMGetLedger messages per peer - Allow a retry after 30s in case of peer or network congestion. - Addresses RIPD-1870 - (Changes levelization. That is not desirable, and will need to be fixed.) - Drop duplicate incoming TMGetLedger messages per peer - Allow a retry after 15s in case of peer or network congestion. - The requestCookie is ignored when computing the hash, thus increasing the chances of detecting duplicate messages. - With duplicate messages, keep track of the different requestCookies (or lack of cookie). When work is finally done for a given request, send the response to all the peers that are waiting on the request, sending one message per peer, including all the cookies and a "directResponse" flag indicating the data is intended for the sender, too. - Addresses RIPD-1871 - Drop duplicate incoming TMLedgerData messages - Addresses RIPD-1869 --- include/xrpl/basics/CanProcess.h | 134 ++++++ include/xrpl/core/HashRouter.h | 38 +- include/xrpl/proto/xrpl.proto | 10 + include/xrpl/protocol/LedgerHeader.h | 2 + include/xrpl/server/NetworkOPs.h | 2 +- src/libxrpl/core/HashRouter.cpp | 22 + src/test/app/HashRouter_test.cpp | 28 ++ src/test/app/LedgerReplay_test.cpp | 5 + src/test/overlay/ProtocolVersion_test.cpp | 4 +- src/test/overlay/reduce_relay_test.cpp | 5 + src/xrpld/app/consensus/RCLConsensus.cpp | 2 +- src/xrpld/app/ledger/InboundLedger.h | 18 + src/xrpld/app/ledger/detail/InboundLedger.cpp | 9 +- .../app/ledger/detail/InboundLedgers.cpp | 122 ++++- src/xrpld/app/ledger/detail/LedgerMaster.cpp | 5 +- .../app/ledger/detail/TimeoutCounter.cpp | 8 +- src/xrpld/app/ledger/detail/TimeoutCounter.h | 2 + src/xrpld/app/misc/NetworkOPs.cpp | 71 ++- src/xrpld/overlay/Peer.h | 8 + src/xrpld/overlay/detail/PeerImp.cpp | 418 ++++++++++++++++-- src/xrpld/overlay/detail/PeerImp.h | 36 +- src/xrpld/overlay/detail/PeerSet.cpp | 35 +- src/xrpld/overlay/detail/ProtocolMessage.h | 66 +++ src/xrpld/overlay/detail/ProtocolVersion.cpp | 2 + 24 files changed, 921 insertions(+), 131 deletions(-) create mode 100644 include/xrpl/basics/CanProcess.h diff --git a/include/xrpl/basics/CanProcess.h b/include/xrpl/basics/CanProcess.h new file mode 100644 index 0000000000..3ee49d0087 --- /dev/null +++ b/include/xrpl/basics/CanProcess.h @@ -0,0 +1,134 @@ +//------------------------------------------------------------------------------ +/* + This file is part of rippled: https://github.com/ripple/rippled + Copyright (c) 2024 Ripple Labs Inc. + + Permission to use, copy, modify, and/or distribute this software for any + purpose with or without fee is hereby granted, provided that the above + copyright notice and this permission notice appear in all copies. + + THE SOFTWARE IS PROVIDED "AS IS" AND THE AUTHOR DISCLAIMS ALL WARRANTIES + WITH REGARD TO THIS SOFTWARE INCLUDING ALL IMPLIED WARRANTIES OF + MERCHANTABILITY AND FITNESS. IN NO EVENT SHALL THE AUTHOR BE LIABLE FOR + ANY SPECIAL , DIRECT, INDIRECT, OR CONSEQUENTIAL DAMAGES OR ANY DAMAGES + WHATSOEVER RESULTING FROM LOSS OF USE, DATA OR PROFITS, WHETHER IN AN + ACTION OF CONTRACT, NEGLIGENCE OR OTHER TORTIOUS ACTION, ARISING OUT OF + OR IN CONNECTION WITH THE USE OR PERFORMANCE OF THIS SOFTWARE. +*/ +//============================================================================== + +#ifndef RIPPLE_BASICS_CANPROCESS_H_INCLUDED +#define RIPPLE_BASICS_CANPROCESS_H_INCLUDED + +#include +#include +#include + +/** RAII class to check if an Item is already being processed on another thread, + * as indicated by it's presence in a Collection. + * + * If the Item is not in the Collection, it will be added under lock in the + * ctor, and removed under lock in the dtor. The object will be considered + * "usable" and evaluate to `true`. + * + * If the Item is in the Collection, no changes will be made to the collection, + * and the CanProcess object will be considered "unusable". + * + * It's up to the caller to decide what "usable" and "unusable" mean. (e.g. + * Process or skip a block of code, or set a flag.) + * + * The current use is to avoid lock contention that would be involved in + * processing something associated with the Item. + * + * Examples: + * + * void IncomingLedgers::acquireAsync(LedgerHash const& hash, ...) + * { + * if (CanProcess check{acquiresMutex_, pendingAcquires_, hash}) + * { + * acquire(hash, ...); + * } + * } + * + * bool + * NetworkOPsImp::recvValidation( + * std::shared_ptr const& val, + * std::string const& source) + * { + * CanProcess check( + * validationsMutex_, pendingValidations_, val->getLedgerHash()); + * BypassAccept bypassAccept = + * check ? BypassAccept::no : BypassAccept::yes; + * handleNewValidation(app_, val, source, bypassAccept, m_journal); + * } + * + */ +class CanProcess +{ +public: + template + CanProcess(Mutex& mtx, Collection& collection, Item const& item) + : cleanup_(insert(mtx, collection, item)) + { + } + + ~CanProcess() + { + if (cleanup_) + cleanup_(); + } + + explicit + operator bool() const + { + return static_cast(cleanup_); + } + +private: + template + std::function + doInsert(Mutex& mtx, Collection& collection, Item const& item) + { + std::unique_lock lock(mtx); + // TODO: Use structured binding once LLVM 16 is the minimum supported + // version. See also: https://github.com/llvm/llvm-project/issues/48582 + // https://github.com/llvm/llvm-project/commit/127bf44385424891eb04cff8e52d3f157fc2cb7c + auto const insertResult = collection.insert(item); + auto const it = insertResult.first; + if (!insertResult.second) + return {}; + if constexpr (useIterator) + return [&, it]() { + std::unique_lock lock(mtx); + collection.erase(it); + }; + else + return [&]() { + std::unique_lock lock(mtx); + collection.erase(item); + }; + } + + // Generic insert() function doesn't use iterators because they may get + // invalidated + template + std::function + insert(Mutex& mtx, Collection& collection, Item const& item) + { + return doInsert(mtx, collection, item); + } + + // Specialize insert() for std::set, which does not invalidate iterators for + // insert and erase + template + std::function + insert(Mutex& mtx, std::set& collection, Item const& item) + { + return doInsert(mtx, collection, item); + } + + // If set, then the item is "usable" + std::function cleanup_; +}; + +#endif diff --git a/include/xrpl/core/HashRouter.h b/include/xrpl/core/HashRouter.h index b4f07f6dc0..de50dd99d4 100644 --- a/include/xrpl/core/HashRouter.h +++ b/include/xrpl/core/HashRouter.h @@ -139,6 +139,13 @@ private: return std::move(peers_); } + /** Return set of peers waiting for reply. Leaves list unchanged. */ + std::set const& + peekPeerSet() + { + return peers_; + } + /** Return seated relay time point if the message has been relayed */ std::optional relayed() const @@ -170,6 +177,20 @@ private: return true; } + bool + shouldProcessForPeer( + PeerShortID peer, + Stopwatch::time_point now, + std::chrono::seconds interval) + { + if (peerProcessed_.contains(peer) && ((peerProcessed_[peer] + interval) > now)) + return false; + // Peer may already be in the list, but adding it again doesn't hurt + addPeer(peer); + peerProcessed_[peer] = now; + return true; + } + private: HashRouterFlags flags_ = HashRouterFlags::UNDEFINED; std::set peers_; @@ -177,6 +198,7 @@ private: // than one flag needs to expire independently. std::optional relayed_; std::optional processed_; + std::map peerProcessed_; }; public: @@ -199,7 +221,7 @@ public: /** Add a suppression peer and get message's relay status. * Return pair: - * element 1: true if the peer is added. + * element 1: true if the key is added. * element 2: optional is seated to the relay time point or * is unseated if has not relayed yet. */ std::pair> @@ -216,6 +238,15 @@ public: HashRouterFlags& flags, std::chrono::seconds tx_interval); + /** Determines whether the hashed item should be processed for the given + peer. Could be an incoming or outgoing message. + + Items filtered with this function should only be processed for the given + peer once. Unlike shouldProcess, it can be processed for other peers. + */ + bool + shouldProcessForPeer(uint256 const& key, PeerShortID peer, std::chrono::seconds interval); + /** Set the flags on a hash. @return `true` if the flags were changed. `false` if unchanged. @@ -241,6 +272,11 @@ public: std::optional> shouldRelay(uint256 const& key); + /** Returns a copy of the set of peers in the Entry for the key + */ + std::set + getPeers(uint256 const& key); + private: // pair.second indicates whether the entry was created std::pair diff --git a/include/xrpl/proto/xrpl.proto b/include/xrpl/proto/xrpl.proto index cd82ed24e6..09d0d2cd25 100644 --- a/include/xrpl/proto/xrpl.proto +++ b/include/xrpl/proto/xrpl.proto @@ -286,8 +286,18 @@ message TMLedgerData { required uint32 ledgerSeq = 2; required TMLedgerInfoType type = 3; repeated TMLedgerNode nodes = 4; + // If the peer supports "responseCookies", this field will + // never be populated. optional uint32 requestCookie = 5; optional TMReplyError error = 6; + // The old field is called "requestCookie", but this is + // a response, so this name makes more sense + repeated uint32 responseCookies = 7; + // If a TMGetLedger request was received without a "requestCookie", + // and the peer supports it, this flag will be set to true to + // indicate that the receiver should process the result in addition + // to forwarding it to its "responseCookies" peers. + optional bool directResponse = 8; } message TMPing { diff --git a/include/xrpl/protocol/LedgerHeader.h b/include/xrpl/protocol/LedgerHeader.h index 6e22ad268d..62050f83fa 100644 --- a/include/xrpl/protocol/LedgerHeader.h +++ b/include/xrpl/protocol/LedgerHeader.h @@ -35,6 +35,8 @@ struct LedgerHeader // If validated is false, it means "not yet validated." // Once validated is true, it will never be set false at a later time. + // NOTE: If you are accessing this directly, you are probably doing it + // wrong. Use LedgerMaster::isValidated(). // VFALCO TODO Make this not mutable bool mutable validated = false; bool accepted = false; diff --git a/include/xrpl/server/NetworkOPs.h b/include/xrpl/server/NetworkOPs.h index 16ec4a4ec0..cf5f6081a3 100644 --- a/include/xrpl/server/NetworkOPs.h +++ b/include/xrpl/server/NetworkOPs.h @@ -185,7 +185,7 @@ public: virtual bool isFull() = 0; virtual void - setMode(OperatingMode om) = 0; + setMode(OperatingMode om, char const* reason) = 0; virtual bool isBlocked() = 0; virtual bool diff --git a/src/libxrpl/core/HashRouter.cpp b/src/libxrpl/core/HashRouter.cpp index f21daf84a2..2dbd60a435 100644 --- a/src/libxrpl/core/HashRouter.cpp +++ b/src/libxrpl/core/HashRouter.cpp @@ -70,6 +70,19 @@ HashRouter::shouldProcess( return s.shouldProcess(suppressionMap_.clock().now(), tx_interval); } +bool +HashRouter::shouldProcessForPeer( + uint256 const& key, + PeerShortID peer, + std::chrono::seconds interval) +{ + std::lock_guard lock(mutex_); + + auto& entry = emplace(key).first; + + return entry.shouldProcessForPeer(peer, suppressionMap_.clock().now(), interval); +} + HashRouterFlags HashRouter::getFlags(uint256 const& key) { @@ -107,4 +120,13 @@ HashRouter::shouldRelay(uint256 const& key) -> std::optional std::set +{ + std::lock_guard lock(mutex_); + + auto& s = emplace(key).first; + return s.peekPeerSet(); +} + } // namespace xrpl diff --git a/src/test/app/HashRouter_test.cpp b/src/test/app/HashRouter_test.cpp index e53515e421..2a94c68816 100644 --- a/src/test/app/HashRouter_test.cpp +++ b/src/test/app/HashRouter_test.cpp @@ -388,6 +388,33 @@ class HashRouter_test : public beast::unit_test::suite BEAST_EXPECT(!any(HF::UNDEFINED)); } + void + testProcessPeer() + { + using namespace std::chrono_literals; + TestStopwatch stopwatch; + HashRouter router(getSetup(5s, 5s), stopwatch); + uint256 const key(1); + HashRouter::PeerShortID peer1 = 1; + HashRouter::PeerShortID peer2 = 2; + auto const timeout = 2s; + + BEAST_EXPECT(router.shouldProcessForPeer(key, peer1, timeout)); + BEAST_EXPECT(!router.shouldProcessForPeer(key, peer1, timeout)); + ++stopwatch; + BEAST_EXPECT(!router.shouldProcessForPeer(key, peer1, timeout)); + BEAST_EXPECT(router.shouldProcessForPeer(key, peer2, timeout)); + BEAST_EXPECT(!router.shouldProcessForPeer(key, peer2, timeout)); + ++stopwatch; + BEAST_EXPECT(router.shouldProcessForPeer(key, peer1, timeout)); + BEAST_EXPECT(!router.shouldProcessForPeer(key, peer2, timeout)); + ++stopwatch; + BEAST_EXPECT(router.shouldProcessForPeer(key, peer2, timeout)); + ++stopwatch; + BEAST_EXPECT(router.shouldProcessForPeer(key, peer1, timeout)); + BEAST_EXPECT(!router.shouldProcessForPeer(key, peer2, timeout)); + } + public: void run() override @@ -400,6 +427,7 @@ public: testProcess(); testSetup(); testFlagsOps(); + testProcessPeer(); } }; diff --git a/src/test/app/LedgerReplay_test.cpp b/src/test/app/LedgerReplay_test.cpp index f9ab08e900..8abf6ee1a0 100644 --- a/src/test/app/LedgerReplay_test.cpp +++ b/src/test/app/LedgerReplay_test.cpp @@ -294,6 +294,11 @@ public: { return false; } + std::set> + releaseRequestCookies(uint256 const& requestHash) override + { + return {}; + } std::string const& fingerprint() const override diff --git a/src/test/overlay/ProtocolVersion_test.cpp b/src/test/overlay/ProtocolVersion_test.cpp index fc25812cbb..a5959eba6f 100644 --- a/src/test/overlay/ProtocolVersion_test.cpp +++ b/src/test/overlay/ProtocolVersion_test.cpp @@ -60,8 +60,8 @@ public: negotiateProtocolVersion("RTXP/1.2, XRPL/2.0, XRPL/2.1") == make_protocol(2, 1)); BEAST_EXPECT(negotiateProtocolVersion("XRPL/2.2") == make_protocol(2, 2)); BEAST_EXPECT( - negotiateProtocolVersion("RTXP/1.2, XRPL/2.2, XRPL/2.3, XRPL/999.999") == - make_protocol(2, 2)); + negotiateProtocolVersion("RTXP/1.2, XRPL/2.2, XRPL/2.3, XRPL/2.4, XRPL/999.999") == + make_protocol(2, 3)); BEAST_EXPECT(negotiateProtocolVersion("XRPL/999.999, WebSocket/1.0") == std::nullopt); BEAST_EXPECT(negotiateProtocolVersion("") == std::nullopt); } diff --git a/src/test/overlay/reduce_relay_test.cpp b/src/test/overlay/reduce_relay_test.cpp index bac70d35a6..bc63fbb923 100644 --- a/src/test/overlay/reduce_relay_test.cpp +++ b/src/test/overlay/reduce_relay_test.cpp @@ -168,6 +168,11 @@ public: removeTxQueue(uint256 const&) override { } + std::set> + releaseRequestCookies(uint256 const& requestHash) override + { + return {}; + } }; /** Manually advanced clock. */ diff --git a/src/xrpld/app/consensus/RCLConsensus.cpp b/src/xrpld/app/consensus/RCLConsensus.cpp index b7b0919aad..83fd80d0c4 100644 --- a/src/xrpld/app/consensus/RCLConsensus.cpp +++ b/src/xrpld/app/consensus/RCLConsensus.cpp @@ -998,7 +998,7 @@ void RCLConsensus::Adaptor::updateOperatingMode(std::size_t const positions) const { if ((positions == 0u) && app_.getOPs().isFull()) - app_.getOPs().setMode(OperatingMode::CONNECTED); + app_.getOPs().setMode(OperatingMode::CONNECTED, "updateOperatingMode: no positions"); } void diff --git a/src/xrpld/app/ledger/InboundLedger.h b/src/xrpld/app/ledger/InboundLedger.h index b17b59b27f..136a0d0ab7 100644 --- a/src/xrpld/app/ledger/InboundLedger.h +++ b/src/xrpld/app/ledger/InboundLedger.h @@ -172,4 +172,22 @@ private: std::unique_ptr mPeerSet; }; +inline std::string +to_string(InboundLedger::Reason reason) +{ + using enum InboundLedger::Reason; + switch (reason) + { + case HISTORY: + return "HISTORY"; + case GENERIC: + return "GENERIC"; + case CONSENSUS: + return "CONSENSUS"; + default: + UNREACHABLE("ripple::to_string(InboundLedger::Reason) : unknown value"); + return "unknown"; + } +} + } // namespace xrpl diff --git a/src/xrpld/app/ledger/detail/InboundLedger.cpp b/src/xrpld/app/ledger/detail/InboundLedger.cpp index 2402b5b561..36c41e1d27 100644 --- a/src/xrpld/app/ledger/detail/InboundLedger.cpp +++ b/src/xrpld/app/ledger/detail/InboundLedger.cpp @@ -353,7 +353,14 @@ InboundLedger::onTimer(bool wasProgress, ScopedLockType&) if (!wasProgress) { - checkLocal(); + if (checkLocal()) + { + // Done. Something else (probably consensus) built the ledger + // locally while waiting for data (or possibly before requesting) + XRPL_ASSERT(isDone(), "ripple::InboundLedger::onTimer : done"); + JLOG(journal_.info()) << "Finished while waiting " << hash_; + return; + } mByHash = true; diff --git a/src/xrpld/app/ledger/detail/InboundLedgers.cpp b/src/xrpld/app/ledger/detail/InboundLedgers.cpp index f147a35ca4..d62a987fb3 100644 --- a/src/xrpld/app/ledger/detail/InboundLedgers.cpp +++ b/src/xrpld/app/ledger/detail/InboundLedgers.cpp @@ -2,6 +2,7 @@ #include #include +#include #include #include #include @@ -53,10 +54,78 @@ public: XRPL_ASSERT( hash.isNonZero(), "xrpl::InboundLedgersImp::acquire::doAcquire : nonzero hash"); - // probably not the right rule - if (app_.getOPs().isNeedNetworkLedger() && (reason != InboundLedger::Reason::GENERIC) && - (reason != InboundLedger::Reason::CONSENSUS)) + bool const needNetworkLedger = app_.getOPs().isNeedNetworkLedger(); + bool const shouldAcquire = [&]() { + if (!needNetworkLedger) + return true; + if (reason == InboundLedger::Reason::GENERIC) + return true; + if (reason == InboundLedger::Reason::CONSENSUS) + return true; + return false; + }(); + + std::stringstream ss; + ss << "InboundLedger::acquire: " + << "Request: " << to_string(hash) << ", " << seq + << " NeedNetworkLedger: " << (needNetworkLedger ? "yes" : "no") + << " Reason: " << to_string(reason) + << " Should acquire: " << (shouldAcquire ? "true." : "false."); + + /* Acquiring ledgers is somewhat expensive. It requires lots of + * computation and network communication. Avoid it when it's not + * appropriate. Every validation from a peer for a ledger that + * we do not have locally results in a call to this function: even + * if we are moments away from validating the same ledger. + */ + bool const shouldBroadcast = [&]() { + // If the node is not in "full" state, it needs to sync to + // the network, and doesn't have the necessary tx's and + // ledger entries to build the ledger. + bool const isFull = app_.getOPs().isFull(); + // If everything else is ok, don't try to acquire the ledger + // if the requested seq is in the near future relative to + // the validated ledger. If the requested ledger is between + // 1 and 19 inclusive ledgers ahead of the valid ledger this + // node has not built it yet, but it's possible/likely it + // has the tx's necessary to build it and get caught up. + // Plus it might not become validated. On the other hand, if + // it's more than 20 in the future, this node should request + // it so that it can jump ahead and get caught up. + LedgerIndex const validSeq = app_.getLedgerMaster().getValidLedgerIndex(); + constexpr std::size_t lagLeeway = 20; + bool const nearFuture = (seq > validSeq) && (seq < validSeq + lagLeeway); + // If everything else is ok, don't try to acquire the ledger + // if the request is related to consensus. (Note that + // consensus calls usually pass a seq of 0, so nearFuture + // will be false other than on a brand new network.) + bool const consensus = reason == InboundLedger::Reason::CONSENSUS; + ss << " Evaluating whether to broadcast requests to peers" + << ". full: " << (isFull ? "true" : "false") << ". ledger sequence " << seq + << ". Valid sequence: " << validSeq << ". Lag leeway: " << lagLeeway + << ". request for near future ledger: " << (nearFuture ? "true" : "false") + << ". Consensus: " << (consensus ? "true" : "false"); + + // If the node is not synced, send requests. + if (!isFull) + return true; + // If the ledger is in the near future, do NOT send requests. + // This node is probably about to build it. + if (nearFuture) + return false; + // If the request is because of consensus, do NOT send requests. + // This node is probably about to build it. + if (consensus) + return false; + return true; + }(); + ss << ". Would broadcast to peers? " << (shouldBroadcast ? "true." : "false."); + + if (!shouldAcquire) + { + JLOG(j_.debug()) << "Abort(rule): " << ss.str(); return {}; + } bool isNew = true; std::shared_ptr inbound; @@ -64,6 +133,7 @@ public: ScopedLockType sl(mLock); if (stopping_) { + JLOG(j_.debug()) << "Abort(stopping): " << ss.str(); return {}; } @@ -82,47 +152,51 @@ public: ++mCounter; } } + ss << " IsNew: " << (isNew ? "true" : "false"); if (inbound->isFailed()) + { + JLOG(j_.debug()) << "Abort(failed): " << ss.str(); return {}; + } if (!isNew) inbound->update(seq); if (!inbound->isComplete()) + { + JLOG(j_.debug()) << "InProgress: " << ss.str(); return {}; + } + JLOG(j_.debug()) << "Complete: " << ss.str(); return inbound->getLedger(); }; using namespace std::chrono_literals; - std::shared_ptr ledger = - perf::measureDurationAndLog(doAcquire, "InboundLedgersImp::acquire", 500ms, j_); - - return ledger; + return perf::measureDurationAndLog(doAcquire, "InboundLedgersImp::acquire", 500ms, j_); } void acquireAsync(uint256 const& hash, std::uint32_t seq, InboundLedger::Reason reason) override { - std::unique_lock lock(acquiresMutex_); - try + if (CanProcess const check{acquiresMutex_, pendingAcquires_, hash}) { - if (pendingAcquires_.contains(hash)) - return; - pendingAcquires_.insert(hash); - scope_unlock const unlock(lock); - acquire(hash, seq, reason); + try + { + acquire(hash, seq, reason); + } + catch (std::exception const& e) + { + JLOG(j_.warn()) << "Exception thrown for acquiring new inbound ledger " << hash + << ": " << e.what(); + } + catch (...) + { + JLOG(j_.warn()) << "Unknown exception thrown for acquiring new " + "inbound ledger " + << hash; + } } - catch (std::exception const& e) - { - JLOG(j_.warn()) << "Exception thrown for acquiring new inbound ledger " << hash << ": " - << e.what(); - } - catch (...) - { - JLOG(j_.warn()) << "Unknown exception thrown for acquiring new inbound ledger " << hash; - } - pendingAcquires_.erase(hash); } std::shared_ptr diff --git a/src/xrpld/app/ledger/detail/LedgerMaster.cpp b/src/xrpld/app/ledger/detail/LedgerMaster.cpp index b9305b743f..d9718b004b 100644 --- a/src/xrpld/app/ledger/detail/LedgerMaster.cpp +++ b/src/xrpld/app/ledger/detail/LedgerMaster.cpp @@ -924,8 +924,9 @@ LedgerMaster::checkAccept(std::shared_ptr const& ledger) return; } - JLOG(m_journal.info()) << "Advancing accepted ledger to " << ledger->header().seq - << " with >= " << minVal << " validations"; + JLOG(m_journal.info()) << "Advancing accepted ledger to " << ledger->header().seq << " (" + << to_short_string(ledger->header().hash) << ") with >= " << minVal + << " validations"; ledger->setValidated(); ledger->setFull(); diff --git a/src/xrpld/app/ledger/detail/TimeoutCounter.cpp b/src/xrpld/app/ledger/detail/TimeoutCounter.cpp index 216771e60d..29d9a228a7 100644 --- a/src/xrpld/app/ledger/detail/TimeoutCounter.cpp +++ b/src/xrpld/app/ledger/detail/TimeoutCounter.cpp @@ -13,7 +13,8 @@ TimeoutCounter::TimeoutCounter( QueueJobParameter&& jobParameter, beast::Journal journal) : app_(app) - , journal_(journal) + , sink_(journal, to_short_string(hash) + " ") + , journal_(sink_) , hash_(hash) , timerInterval_(interval) , queueJobParameter_(std::move(jobParameter)) @@ -29,6 +30,7 @@ TimeoutCounter::setTimer(ScopedLockType& sl) { if (isDone()) return; + JLOG(journal_.debug()) << "Setting timer for " << timerInterval_.count() << "ms"; timer_.expires_after(timerInterval_); timer_.async_wait([wptr = pmDowncast()](boost::system::error_code const& ec) { if (ec == boost::asio::error::operation_aborted) @@ -36,6 +38,10 @@ TimeoutCounter::setTimer(ScopedLockType& sl) if (auto ptr = wptr.lock()) { + JLOG(ptr->journal_.debug()) + << "timer: ec: " << ec + << " (operation_aborted: " << boost::asio::error::operation_aborted << " - " + << (ec == boost::asio::error::operation_aborted ? "aborted" : "other") << ")"; ScopedLockType sl(ptr->mtx_); ptr->queueJob(sl); } diff --git a/src/xrpld/app/ledger/detail/TimeoutCounter.h b/src/xrpld/app/ledger/detail/TimeoutCounter.h index a7e4c043be..055932cbee 100644 --- a/src/xrpld/app/ledger/detail/TimeoutCounter.h +++ b/src/xrpld/app/ledger/detail/TimeoutCounter.h @@ -3,6 +3,7 @@ #include #include +#include #include #include @@ -103,6 +104,7 @@ protected: // Used in this class for access to boost::asio::io_context and // xrpl::Overlay. Used in subtypes for the kitchen sink. Application& app_; + beast::WrappedSink sink_; beast::Journal journal_; mutable std::recursive_mutex mtx_; diff --git a/src/xrpld/app/misc/NetworkOPs.cpp b/src/xrpld/app/misc/NetworkOPs.cpp index e39230efdb..705d08f13b 100644 --- a/src/xrpld/app/misc/NetworkOPs.cpp +++ b/src/xrpld/app/misc/NetworkOPs.cpp @@ -30,10 +30,10 @@ #include #include +#include #include #include #include -#include #include #include #include @@ -401,7 +401,7 @@ public: isFull() override; void - setMode(OperatingMode om) override; + setMode(OperatingMode om, char const* reason) override; bool isBlocked() override; @@ -839,7 +839,7 @@ NetworkOPsImp::strOperatingMode(bool const admin /* = false */) const inline void NetworkOPsImp::setStandAlone() { - setMode(OperatingMode::FULL); + setMode(OperatingMode::FULL, "setStandAlone"); } inline void @@ -982,7 +982,7 @@ NetworkOPsImp::processHeartbeatTimer() { if (mMode != OperatingMode::DISCONNECTED) { - setMode(OperatingMode::DISCONNECTED); + setMode(OperatingMode::DISCONNECTED, "Heartbeat: insufficient peers"); std::stringstream ss; ss << "Node count (" << numPeers << ") has fallen " << "below required minimum (" << minPeerCount_ << ")."; @@ -1006,7 +1006,7 @@ NetworkOPsImp::processHeartbeatTimer() if (mMode == OperatingMode::DISCONNECTED) { - setMode(OperatingMode::CONNECTED); + setMode(OperatingMode::CONNECTED, "Heartbeat: sufficient peers"); JLOG(m_journal.info()) << "Node count (" << numPeers << ") is sufficient."; CLOG(clog.ss()) << "setting mode to CONNECTED based on " << numPeers << " peers. "; } @@ -1017,11 +1017,11 @@ NetworkOPsImp::processHeartbeatTimer() CLOG(clog.ss()) << "mode: " << strOperatingMode(origMode, true); if (mMode == OperatingMode::SYNCING) { - setMode(OperatingMode::SYNCING); + setMode(OperatingMode::SYNCING, "Heartbeat: check syncing"); } else if (mMode == OperatingMode::CONNECTED) { - setMode(OperatingMode::CONNECTED); + setMode(OperatingMode::CONNECTED, "Heartbeat: check connected"); } auto newMode = mMode.load(); if (origMode != newMode) @@ -1726,7 +1726,7 @@ void NetworkOPsImp::setAmendmentBlocked() { amendmentBlocked_ = true; - setMode(OperatingMode::CONNECTED); + setMode(OperatingMode::CONNECTED, "setAmendmentBlocked"); } inline bool @@ -1757,7 +1757,7 @@ void NetworkOPsImp::setUNLBlocked() { unlBlocked_ = true; - setMode(OperatingMode::CONNECTED); + setMode(OperatingMode::CONNECTED, "setUNLBlocked"); } inline void @@ -1857,7 +1857,7 @@ NetworkOPsImp::checkLastClosedLedger(Overlay::PeerSequence const& peerList, uint if ((mMode == OperatingMode::TRACKING) || (mMode == OperatingMode::FULL)) { - setMode(OperatingMode::CONNECTED); + setMode(OperatingMode::CONNECTED, "check LCL: not on consensus ledger"); } if (consensus) @@ -1945,8 +1945,8 @@ NetworkOPsImp::beginConsensus( // this shouldn't happen unless we jump ledgers if (mMode == OperatingMode::FULL) { - JLOG(m_journal.warn()) << "Don't have LCL, going to tracking"; - setMode(OperatingMode::TRACKING); + JLOG(m_journal.warn()) << "beginConsensus Don't have LCL, going to tracking"; + setMode(OperatingMode::TRACKING, "beginConsensus: No LCL"); CLOG(clog) << "beginConsensus Don't have LCL, going to tracking. "; } @@ -2074,7 +2074,7 @@ NetworkOPsImp::endConsensus(std::unique_ptr const& clog) // validations we have for LCL. If the ledger is good enough, go to // TRACKING - TODO if (!needNetworkLedger_) - setMode(OperatingMode::TRACKING); + setMode(OperatingMode::TRACKING, "endConsensus: check tracking"); } if (((mMode == OperatingMode::CONNECTED) || (mMode == OperatingMode::TRACKING)) && @@ -2087,7 +2087,7 @@ NetworkOPsImp::endConsensus(std::unique_ptr const& clog) if (registry_.get().getTimeKeeper().now() < (current->header().parentCloseTime + 2 * current->header().closeTimeResolution)) { - setMode(OperatingMode::FULL); + setMode(OperatingMode::FULL, "endConsensus: check full"); } } @@ -2099,7 +2099,7 @@ NetworkOPsImp::consensusViewChange() { if ((mMode == OperatingMode::FULL) || (mMode == OperatingMode::TRACKING)) { - setMode(OperatingMode::CONNECTED); + setMode(OperatingMode::CONNECTED, "consensusViewChange"); } } @@ -2403,7 +2403,7 @@ NetworkOPsImp::pubPeerStatus(std::function const& func) } void -NetworkOPsImp::setMode(OperatingMode om) +NetworkOPsImp::setMode(OperatingMode om, char const* reason) { using namespace std::chrono_literals; if (om == OperatingMode::CONNECTED) @@ -2423,11 +2423,12 @@ NetworkOPsImp::setMode(OperatingMode om) if (mMode == om) return; + auto const sink = om < mMode ? m_journal.warn() : m_journal.info(); mMode = om; accounting_.mode(om); - JLOG(m_journal.info()) << "STATE->" << strOperatingMode(); + JLOG(sink) << "STATE->" << strOperatingMode() << " - " << reason; pubServer(); } @@ -2436,36 +2437,24 @@ NetworkOPsImp::recvValidation(std::shared_ptr const& val, std::str { JLOG(m_journal.trace()) << "recvValidation " << val->getLedgerHash() << " from " << source; - std::unique_lock lock(validationsMutex_); - BypassAccept bypassAccept = BypassAccept::no; - try { - if (pendingValidations_.contains(val->getLedgerHash())) + CanProcess const check(validationsMutex_, pendingValidations_, val->getLedgerHash()); + try { - bypassAccept = BypassAccept::yes; + BypassAccept bypassAccept = check ? BypassAccept::no : BypassAccept::yes; + handleNewValidation(registry_.app(), val, source, bypassAccept, m_journal); } - else + catch (std::exception const& e) { - pendingValidations_.insert(val->getLedgerHash()); + JLOG(m_journal.warn()) << "Exception thrown for handling new validation " + << val->getLedgerHash() << ": " << e.what(); + } + catch (...) + { + JLOG(m_journal.warn()) + << "Unknown exception thrown for handling new validation " << val->getLedgerHash(); } - scope_unlock const unlock(lock); - handleNewValidation(registry_.get().getApp(), val, source, bypassAccept, m_journal); } - catch (std::exception const& e) - { - JLOG(m_journal.warn()) << "Exception thrown for handling new validation " - << val->getLedgerHash() << ": " << e.what(); - } - catch (...) - { - JLOG(m_journal.warn()) << "Unknown exception thrown for handling new validation " - << val->getLedgerHash(); - } - if (bypassAccept == BypassAccept::no) - { - pendingValidations_.erase(val->getLedgerHash()); - } - lock.unlock(); pubValidation(val); diff --git a/src/xrpld/overlay/Peer.h b/src/xrpld/overlay/Peer.h index df2cc5bcb7..aa145ce9a8 100644 --- a/src/xrpld/overlay/Peer.h +++ b/src/xrpld/overlay/Peer.h @@ -17,6 +17,7 @@ enum class ProtocolFeature { ValidatorListPropagation, ValidatorList2Propagation, LedgerReplay, + LedgerDataCookies }; /** Represents a peer connection in the overlay. */ @@ -116,6 +117,13 @@ public: virtual bool txReduceRelayEnabled() const = 0; + + // + // Messages + // + + virtual std::set> + releaseRequestCookies(uint256 const& requestHash) = 0; }; } // namespace xrpl diff --git a/src/xrpld/overlay/detail/PeerImp.cpp b/src/xrpld/overlay/detail/PeerImp.cpp index 13fe0c571c..e519b64ecc 100644 --- a/src/xrpld/overlay/detail/PeerImp.cpp +++ b/src/xrpld/overlay/detail/PeerImp.cpp @@ -7,6 +7,7 @@ #include #include #include +#include #include #include @@ -45,6 +46,8 @@ std::chrono::seconds constexpr peerTimerInterval{60}; /** The timeout for a shutdown timer */ std::chrono::seconds constexpr shutdownTimerInterval{5}; +/** How often we process duplicate incoming TMGetLedger messages */ +std::chrono::seconds constexpr getledgerInterval{15}; } // namespace // TODO: Remove this exclusion once unit tests are added after the hotfix @@ -500,6 +503,8 @@ PeerImp::supportsFeature(ProtocolFeature f) const return protocol_ >= make_protocol(2, 2); case ProtocolFeature::LedgerReplay: return ledgerReplayEnabled_; + case ProtocolFeature::LedgerDataCookies: + return protocol_ >= make_protocol(2, 3); } return false; } @@ -1481,8 +1486,9 @@ PeerImp::handleTransaction( void PeerImp::onMessage(std::shared_ptr const& m) { - auto badData = [&](std::string const& msg) { - fee_.update(Resource::feeInvalidData, "get_ledger " + msg); + auto badData = [&](std::string const& msg, bool chargefee = true) { + if (chargefee) + fee_.update(Resource::feeInvalidData, "get_ledger " + msg); JLOG(p_journal_.warn()) << "TMGetLedger: " << msg; }; auto const itype{m->itype()}; @@ -1580,11 +1586,68 @@ PeerImp::onMessage(std::shared_ptr const& m) } } + // Drop duplicate requests from the same peer for at least + // `getLedgerInterval` seconds. + // Append a little junk to prevent the hash of an incoming messsage + // from matching the hash of the same outgoing message. + // `shouldProcessForPeer` does not distingish between incoming and + // outgoing, and some of the message relay logic checks the hash to see + // if the message has been relayed already. If the hashes are the same, + // a duplicate will be detected when sending the message is attempted, + // so it will fail. + auto const messageHash = sha512Half(*m, nullptr); + // Request cookies are not included in the hash. Track them here. + auto const requestCookie = [&m]() -> std::optional { + if (m->has_requestcookie()) + return m->requestcookie(); + return std::nullopt; + }(); + auto const [inserted, pending] = [&] { + std::lock_guard lock{cookieLock_}; + auto& cookies = messageRequestCookies_[messageHash]; + bool const pending = !cookies.empty(); + return std::pair{cookies.emplace(requestCookie).second, pending}; + }(); + // Check if the request has been seen from this peer. + if (!app_.getHashRouter().shouldProcessForPeer(messageHash, id_, getledgerInterval)) + { + // This request has already been seen from this peer. + // Has it been seen with this request cookie (or lack thereof)? + + if (inserted) + { + // This is a duplicate request, but with a new cookie. When a + // response is ready, one will be sent for each request cookie. + JLOG(p_journal_.debug()) + << "TMGetLedger: duplicate request with new request cookie: " + << requestCookie.value_or(0) << ". Job pending: " << (pending ? "yes" : "no") + << ": " << messageHash; + if (pending) + { + // Don't bother queueing up a new job if other requests are + // already pending. This should limit entries in the job queue + // to one per peer per unique request. + JLOG(p_journal_.debug()) << "TMGetLedger: Suppressing recvGetLedger job, since one " + "is pending: " + << messageHash; + return; + } + } + else + { + // Don't punish nodes that don't know any better + return badData( + "duplicate request: " + to_string(messageHash), + supportsFeature(ProtocolFeature::LedgerDataCookies)); + } + } + // Queue a job to process the request + JLOG(p_journal_.debug()) << "TMGetLedger: Adding recvGetLedger job: " << messageHash; std::weak_ptr const weak = shared_from_this(); - app_.getJobQueue().addJob(jtLEDGER_REQ, "RcvGetLedger", [weak, m]() { + app_.getJobQueue().addJob(jtLEDGER_REQ, "RcvGetLedger", [weak, m, messageHash]() { if (auto peer = weak.lock()) - peer->processLedgerRequest(m); + peer->processLedgerRequest(m, messageHash); }); } @@ -1691,8 +1754,9 @@ PeerImp::onMessage(std::shared_ptr const& m) void PeerImp::onMessage(std::shared_ptr const& m) { - auto badData = [&](std::string const& msg) { - fee_.update(Resource::feeInvalidData, msg); + auto badData = [&](std::string const& msg, bool charge = true) { + if (charge) + fee_.update(Resource::feeInvalidData, msg); JLOG(p_journal_.warn()) << "TMLedgerData: " << msg; }; @@ -1749,23 +1813,96 @@ PeerImp::onMessage(std::shared_ptr const& m) return; } - // If there is a request cookie, attempt to relay the message - if (m->has_requestcookie()) + auto const messageHash = sha512Half(*m); + if (!app_.getHashRouter().addSuppressionPeer(messageHash, id_)) { - if (auto peer = overlay_.findPeerByShortID(m->requestcookie())) + // Don't punish nodes that don't know any better + return badData( + "Duplicate message: " + to_string(messageHash), + supportsFeature(ProtocolFeature::LedgerDataCookies)); + } + + bool const routed = + m->has_directresponse() || m->responsecookies_size() || m->has_requestcookie(); + + { + // Check if this message needs to be forwarded to one or more peers. + // Maximum of one of the relevant fields should be populated. + XRPL_ASSERT( + !m->has_requestcookie() || !m->responsecookies_size(), + "ripple::PeerImp::onMessage(TMLedgerData) : valid cookie fields"); + + // Make a copy of the response cookies, then wipe the list so it can be + // forwarded cleanly + auto const responseCookies = m->responsecookies(); + m->clear_responsecookies(); + // Flag indicating if this response should be processed locally, + // possibly in addition to being forwarded. + bool const directResponse = m->has_directresponse() && m->directresponse(); + m->clear_directresponse(); + + auto const relay = [this, m, &messageHash](auto const cookie) { + if (auto peer = overlay_.findPeerByShortID(cookie)) + { + XRPL_ASSERT( + !m->has_requestcookie() && !m->responsecookies_size(), + "ripple::PeerImp::onMessage(TMLedgerData) relay : no " + "cookies"); + if (peer->supportsFeature(ProtocolFeature::LedgerDataCookies)) + // Setting this flag is not _strictly_ necessary for peers + // that support it if there are no cookies included in the + // message, but it is more accurate. + m->set_directresponse(true); + else + m->clear_directresponse(); + peer->send(std::make_shared(*m, protocol::mtLEDGER_DATA)); + } + else + JLOG(p_journal_.info()) << "Unable to route TX/ledger data reply to peer [" + << cookie << "]: " << messageHash; + }; + // If there is a request cookie, attempt to relay the message + if (m->has_requestcookie()) { + XRPL_ASSERT( + responseCookies.empty(), + "ripple::PeerImp::onMessage(TMLedgerData) : no response " + "cookies"); m->clear_requestcookie(); - peer->send(std::make_shared(*m, protocol::mtLEDGER_DATA)); + relay(m->requestcookie()); + if (!directResponse && responseCookies.empty()) + return; } - else + // If there's a list of request cookies, attempt to relay the message to + // all of them. + if (responseCookies.size()) { - JLOG(p_journal_.info()) << "Unable to route TX/ledger data reply"; + for (auto const cookie : responseCookies) + relay(cookie); + if (!directResponse) + return; + } + } + + // Now that any forwarding is done check the base message (data only, no + // routing info for duplicates) + if (routed) + { + m->clear_directresponse(); + XRPL_ASSERT( + !m->has_requestcookie() && !m->responsecookies_size(), + "ripple::PeerImp::onMessage(TMLedgerData) : no cookies"); + auto const baseMessageHash = sha512Half(*m); + if (!app_.getHashRouter().addSuppressionPeer(baseMessageHash, id_)) + { + // Don't punish nodes that don't know any better + return badData( + "Duplicate message: " + to_string(baseMessageHash), + supportsFeature(ProtocolFeature::LedgerDataCookies)); } - return; } uint256 const ledgerHash{m->ledgerhash()}; - // Otherwise check if received data for a candidate transaction set if (m->type() == protocol::liTS_CANDIDATE) { @@ -3096,16 +3233,21 @@ PeerImp::checkValidation( // the TX tree with the specified root hash. // static std::shared_ptr -getPeerWithTree(OverlayImpl& ov, uint256 const& rootHash, PeerImp const* skip) +getPeerWithTree( + OverlayImpl& ov, + uint256 const& rootHash, + PeerImp const* skip, + std::function shouldProcessCallback) { std::shared_ptr ret; int retScore = 0; + XRPL_ASSERT(shouldProcessCallback, "ripple::getPeerWithTree : callback provided"); ov.for_each([&](std::shared_ptr&& p) { if (p->hasTxSet(rootHash) && p.get() != skip) { auto score = p->getScore(true); - if (!ret || (score > retScore)) + if (!ret || (score > retScore && shouldProcessCallback(p->id()))) { ret = std::move(p); retScore = score; @@ -3124,16 +3266,18 @@ getPeerWithLedger( OverlayImpl& ov, uint256 const& ledgerHash, LedgerIndex ledger, - PeerImp const* skip) + PeerImp const* skip, + std::function shouldProcessCallback) { std::shared_ptr ret; int retScore = 0; + XRPL_ASSERT(shouldProcessCallback, "ripple::getPeerWithLedger : callback provided"); ov.for_each([&](std::shared_ptr&& p) { if (p->hasLedger(ledgerHash, ledger) && p.get() != skip) { auto score = p->getScore(true); - if (!ret || (score > retScore)) + if (!ret || (score > retScore && shouldProcessCallback(p->id()))) { ret = std::move(p); retScore = score; @@ -3147,7 +3291,8 @@ getPeerWithLedger( void PeerImp::sendLedgerBase( std::shared_ptr const& ledger, - protocol::TMLedgerData& ledgerData) + protocol::TMLedgerData& ledgerData, + PeerCookieMap const& destinations) { JLOG(p_journal_.trace()) << "sendLedgerBase: Base data"; @@ -3177,14 +3322,92 @@ PeerImp::sendLedgerBase( } } - auto message{std::make_shared(ledgerData, protocol::mtLEDGER_DATA)}; - send(message); + sendToMultiple(ledgerData, destinations); +} + +void +PeerImp::sendToMultiple(protocol::TMLedgerData& ledgerData, PeerCookieMap const& destinations) +{ + bool foundSelf = false; + for (auto const& [peer, cookies] : destinations) + { + if (peer.get() == this) + foundSelf = true; + bool const multipleCookies = peer->supportsFeature(ProtocolFeature::LedgerDataCookies); + std::vector sendCookies; + + bool directResponse = false; + if (!multipleCookies) + { + JLOG(p_journal_.debug()) << "sendToMultiple: Sending " << cookies.size() + << " TMLedgerData messages to peer [" << peer->id() + << "]: " << sha512Half(ledgerData); + } + for (auto const& cookie : cookies) + { + // Unfortunately, need a separate Message object for every + // combination + if (cookie) + { + if (multipleCookies) + { + // Save this one for later to send a single message + sendCookies.emplace_back(*cookie); + continue; + } + + // Feature not supported, so send a single message with a + // single cookie + ledgerData.set_requestcookie(*cookie); + } + else + { + if (multipleCookies) + { + // Set this flag later on the single message + directResponse = true; + continue; + } + + ledgerData.clear_requestcookie(); + } + XRPL_ASSERT( + !multipleCookies, + "ripple::PeerImp::sendToMultiple : ledger data cookies " + "unsupported"); + auto message{std::make_shared(ledgerData, protocol::mtLEDGER_DATA)}; + peer->send(message); + } + if (multipleCookies) + { + // Send a single message with all the cookies and/or the direct + // response flag, so the receiver can farm out the single message to + // multiple peers and/or itself + XRPL_ASSERT( + sendCookies.size() || directResponse, + "ripple::PeerImp::sendToMultiple : valid response options"); + ledgerData.clear_requestcookie(); + ledgerData.clear_responsecookies(); + ledgerData.set_directresponse(directResponse); + for (auto const& cookie : sendCookies) + ledgerData.add_responsecookies(cookie); + auto message{std::make_shared(ledgerData, protocol::mtLEDGER_DATA)}; + peer->send(message); + + JLOG(p_journal_.debug()) + << "sendToMultiple: Sent 1 TMLedgerData message to peer [" << peer->id() + << "]: including " << (directResponse ? "the direct response flag and " : "") + << sendCookies.size() << " response cookies. " + << ": " << sha512Half(ledgerData); + } + } + XRPL_ASSERT(foundSelf, "ripple::PeerImp::sendToMultiple : current peer included"); } std::shared_ptr -PeerImp::getLedger(std::shared_ptr const& m) +PeerImp::getLedger(std::shared_ptr const& m, uint256 const& mHash) { - JLOG(p_journal_.trace()) << "getLedger: Ledger"; + JLOG(p_journal_.trace()) << "getLedger: Ledger " << mHash; std::shared_ptr ledger; @@ -3200,16 +3423,30 @@ PeerImp::getLedger(std::shared_ptr const& m) if (m->has_querytype() && !m->has_requestcookie()) { // Attempt to relay the request to a peer + // Note repeated messages will not relay to the same peer + // before `getLedgerInterval` seconds. This prevents one + // peer from getting flooded, and distributes the request + // load. If a request has been relayed to all eligible + // peers, then this message will not be relayed. if (auto const peer = getPeerWithLedger( - overlay_, ledgerHash, m->has_ledgerseq() ? m->ledgerseq() : 0, this)) + overlay_, + ledgerHash, + m->has_ledgerseq() ? m->ledgerseq() : 0, + this, + [&](Peer::id_t id) { + return app_.getHashRouter().shouldProcessForPeer( + mHash, id, getledgerInterval); + })) { m->set_requestcookie(id()); peer->send(std::make_shared(*m, protocol::mtGET_LEDGER)); - JLOG(p_journal_.debug()) << "getLedger: Request relayed to peer"; + JLOG(p_journal_.debug()) + << "getLedger: Request relayed to peer [" << peer->id() << "]: " << mHash; return ledger; } - JLOG(p_journal_.trace()) << "getLedger: Failed to find peer to relay request"; + JLOG(p_journal_.trace()) + << "getLedger: Don't have ledger with hash " << ledgerHash << ": " << mHash; } } } @@ -3218,15 +3455,15 @@ PeerImp::getLedger(std::shared_ptr const& m) // Attempt to find ledger by sequence if (m->ledgerseq() < app_.getLedgerMaster().getEarliestFetch()) { - JLOG(p_journal_.debug()) << "getLedger: Early ledger sequence request"; + JLOG(p_journal_.debug()) << "getLedger: Early ledger sequence request " << mHash; } else { ledger = app_.getLedgerMaster().getLedgerBySeq(m->ledgerseq()); if (!ledger) { - JLOG(p_journal_.debug()) - << "getLedger: Don't have ledger with sequence " << m->ledgerseq(); + JLOG(p_journal_.debug()) << "getLedger: Don't have ledger with sequence " + << m->ledgerseq() << ": " << mHash; } } } @@ -3248,27 +3485,29 @@ PeerImp::getLedger(std::shared_ptr const& m) charge(Resource::feeMalformedRequest, "get_ledger ledgerSeq"); ledger.reset(); - JLOG(p_journal_.warn()) << "getLedger: Invalid ledger sequence " << ledgerSeq; + JLOG(p_journal_.warn()) + << "getLedger: Invalid ledger sequence " << ledgerSeq << ": " << mHash; } } else if (ledgerSeq < app_.getLedgerMaster().getEarliestFetch()) { ledger.reset(); - JLOG(p_journal_.debug()) << "getLedger: Early ledger sequence request " << ledgerSeq; + JLOG(p_journal_.debug()) + << "getLedger: Early ledger sequence request " << ledgerSeq << ": " << mHash; } } else { - JLOG(p_journal_.debug()) << "getLedger: Unable to find ledger"; + JLOG(p_journal_.debug()) << "getLedger: Unable to find ledger " << mHash; } return ledger; } std::shared_ptr -PeerImp::getTxSet(std::shared_ptr const& m) const +PeerImp::getTxSet(std::shared_ptr const& m, uint256 const& mHash) const { - JLOG(p_journal_.trace()) << "getTxSet: TX set"; + JLOG(p_journal_.trace()) << "getTxSet: TX set " << mHash; uint256 const txSetHash{m->ledgerhash()}; std::shared_ptr shaMap{app_.getInboundTransactions().getSet(txSetHash, false)}; @@ -3277,20 +3516,28 @@ PeerImp::getTxSet(std::shared_ptr const& m) const if (m->has_querytype() && !m->has_requestcookie()) { // Attempt to relay the request to a peer - if (auto const peer = getPeerWithTree(overlay_, txSetHash, this)) + // Note repeated messages will not relay to the same peer + // before `getLedgerInterval` seconds. This prevents one + // peer from getting flooded, and distributes the request + // load. If a request has been relayed to all eligible + // peers, then this message will not be relayed. + if (auto const peer = getPeerWithTree(overlay_, txSetHash, this, [&](Peer::id_t id) { + return app_.getHashRouter().shouldProcessForPeer(mHash, id, getledgerInterval); + })) { m->set_requestcookie(id()); peer->send(std::make_shared(*m, protocol::mtGET_LEDGER)); - JLOG(p_journal_.debug()) << "getTxSet: Request relayed"; + JLOG(p_journal_.debug()) + << "getTxSet: Request relayed to peer [" << peer->id() << "]: " << mHash; } else { - JLOG(p_journal_.debug()) << "getTxSet: Failed to find relay peer"; + JLOG(p_journal_.debug()) << "getTxSet: Failed to find relay peer: " << mHash; } } else { - JLOG(p_journal_.debug()) << "getTxSet: Failed to find TX set"; + JLOG(p_journal_.debug()) << "getTxSet: Failed to find TX set " << mHash; } } @@ -3298,7 +3545,7 @@ PeerImp::getTxSet(std::shared_ptr const& m) const } void -PeerImp::processLedgerRequest(std::shared_ptr const& m) +PeerImp::processLedgerRequest(std::shared_ptr const& m, uint256 const& mHash) { // Do not resource charge a peer responding to a relay if (!m->has_requestcookie()) @@ -3311,9 +3558,72 @@ PeerImp::processLedgerRequest(std::shared_ptr const& m) bool fatLeaves{true}; auto const itype{m->itype()}; + auto getDestinations = [&] { + // If a ledger data message is generated, it's going to be sent to every + // peer that is waiting for it. + + PeerCookieMap result; + + std::size_t numCookies = 0; + { + // Don't do the work under this peer if this peer is not waiting for + // any replies + auto myCookies = releaseRequestCookies(mHash); + if (myCookies.empty()) + { + JLOG(p_journal_.debug()) << "TMGetLedger: peer is no longer " + "waiting for response to request: " + << mHash; + return result; + } + numCookies += myCookies.size(); + result[shared_from_this()] = myCookies; + } + + std::set const peers = app_.getHashRouter().getPeers(mHash); + for (auto const peerID : peers) + { + // This loop does not need to be done under the HashRouter + // lock because findPeerByShortID and releaseRequestCookies + // are thread safe, and everything else is local + if (auto p = overlay_.findPeerByShortID(peerID)) + { + auto cookies = p->releaseRequestCookies(mHash); + numCookies += cookies.size(); + if (result.contains(p)) + { + // Unlikely, but if a request came in to this peer while + // iterating, add the items instead of copying / + // overwriting. + XRPL_ASSERT( + p.get() == this, + "ripple::PeerImp::processLedgerRequest : found self in " + "map"); + for (auto const& cookie : cookies) + result[p].emplace(cookie); + } + else if (cookies.size()) + result[p] = cookies; + } + } + + JLOG(p_journal_.debug()) << "TMGetLedger: Processing request for " << result.size() + << " peers. Will send " << numCookies + << " messages if successful: " << mHash; + + return result; + }; + // Will only populate this if we're going to do work. + PeerCookieMap destinations; + if (itype == protocol::liTS_CANDIDATE) { - if (sharedMap = getTxSet(m); !sharedMap) + destinations = getDestinations(); + if (destinations.empty()) + // Nowhere to send the response! + return; + + if (sharedMap = getTxSet(m, mHash); !sharedMap) return; map = sharedMap.get(); @@ -3321,8 +3631,6 @@ PeerImp::processLedgerRequest(std::shared_ptr const& m) ledgerData.set_ledgerseq(0); ledgerData.set_ledgerhash(m->ledgerhash()); ledgerData.set_type(protocol::liTS_CANDIDATE); - if (m->has_requestcookie()) - ledgerData.set_requestcookie(m->requestcookie()); // We'll already have most transactions fatLeaves = false; @@ -3340,7 +3648,12 @@ PeerImp::processLedgerRequest(std::shared_ptr const& m) return; } - if (ledger = getLedger(m); !ledger) + destinations = getDestinations(); + if (destinations.empty()) + // Nowhere to send the response! + return; + + if (ledger = getLedger(m, mHash); !ledger) return; // Fill out the reply @@ -3348,13 +3661,11 @@ PeerImp::processLedgerRequest(std::shared_ptr const& m) ledgerData.set_ledgerhash(ledgerHash.begin(), ledgerHash.size()); ledgerData.set_ledgerseq(ledger->header().seq); ledgerData.set_type(itype); - if (m->has_requestcookie()) - ledgerData.set_requestcookie(m->requestcookie()); switch (itype) { case protocol::liBASE: - sendLedgerBase(ledger, ledgerData); + sendLedgerBase(ledger, ledgerData, destinations); return; case protocol::liTX_NODE: @@ -3464,7 +3775,7 @@ PeerImp::processLedgerRequest(std::shared_ptr const& m) if (ledgerData.nodes_size() == 0) return; - send(std::make_shared(ledgerData, protocol::mtLEDGER_DATA)); + sendToMultiple(ledgerData, destinations); } int @@ -3516,6 +3827,19 @@ PeerImp::isHighLatency() const return latency_ >= peerHighLatency; } +std::set> +PeerImp::releaseRequestCookies(uint256 const& requestHash) +{ + std::set> result; + std::lock_guard lock(cookieLock_); + if (messageRequestCookies_.contains(requestHash)) + { + std::swap(result, messageRequestCookies_[requestHash]); + messageRequestCookies_.erase(requestHash); + } + return result; +}; + void PeerImp::Metrics::add_message(std::uint64_t bytes) { diff --git a/src/xrpld/overlay/detail/PeerImp.h b/src/xrpld/overlay/detail/PeerImp.h index 7f7d8c9324..53713039a2 100644 --- a/src/xrpld/overlay/detail/PeerImp.h +++ b/src/xrpld/overlay/detail/PeerImp.h @@ -246,6 +246,13 @@ private: bool ledgerReplayEnabled_ = false; LedgerReplayMsgHandler ledgerReplayMsgHandler_; + // Track message requests and responses + // TODO: Use an expiring cache or something + using MessageCookieMap = std::map>>; + using PeerCookieMap = std::map, std::set>>; + std::mutex mutable cookieLock_; + MessageCookieMap messageRequestCookies_; + friend class OverlayImpl; class Metrics @@ -488,6 +495,13 @@ public: return txReduceRelayEnabled_; } + // + // Messages + // + + std::set> + releaseRequestCookies(uint256 const& requestHash) override; + private: /** * @brief Handles a failure associated with a specific error code. @@ -783,16 +797,22 @@ private: std::shared_ptr const& packet); void - sendLedgerBase(std::shared_ptr const& ledger, protocol::TMLedgerData& ledgerData); - - std::shared_ptr - getLedger(std::shared_ptr const& m); - - std::shared_ptr - getTxSet(std::shared_ptr const& m) const; + sendLedgerBase( + std::shared_ptr const& ledger, + protocol::TMLedgerData& ledgerData, + PeerCookieMap const& destinations); void - processLedgerRequest(std::shared_ptr const& m); + sendToMultiple(protocol::TMLedgerData& ledgerData, PeerCookieMap const& destinations); + + std::shared_ptr + getLedger(std::shared_ptr const& m, uint256 const& mHash); + + std::shared_ptr + getTxSet(std::shared_ptr const& m, uint256 const& mHash) const; + + void + processLedgerRequest(std::shared_ptr const& m, uint256 const& mHash); }; //------------------------------------------------------------------------------ diff --git a/src/xrpld/overlay/detail/PeerSet.cpp b/src/xrpld/overlay/detail/PeerSet.cpp index 391fb6d3ca..2a980e4ec4 100644 --- a/src/xrpld/overlay/detail/PeerSet.cpp +++ b/src/xrpld/overlay/detail/PeerSet.cpp @@ -2,7 +2,9 @@ #include #include +#include #include +#include namespace xrpl { @@ -82,16 +84,45 @@ PeerSetImpl::sendRequest( std::shared_ptr const& peer) { auto packet = std::make_shared(message, type); + + auto const messageHash = [&]() { + auto const packetBuffer = packet->getBuffer(compression::Compressed::Off); + return sha512Half(Slice(packetBuffer.data(), packetBuffer.size())); + }(); + + // Allow messages to be re-sent to the same peer after a delay + using namespace std::chrono_literals; + constexpr std::chrono::seconds interval = 30s; + if (peer) { - peer->send(packet); + if (app_.getHashRouter().shouldProcessForPeer(messageHash, peer->id(), interval)) + { + JLOG(journal_.trace()) << "Sending " << protocolMessageName(type) << " message to [" + << peer->id() << "]: " << messageHash; + peer->send(packet); + } + else + JLOG(journal_.debug()) << "Suppressing sending duplicate " << protocolMessageName(type) + << " message to [" << peer->id() << "]: " << messageHash; return; } for (auto id : peers_) { if (auto p = app_.getOverlay().findPeerByShortID(id)) - p->send(packet); + { + if (app_.getHashRouter().shouldProcessForPeer(messageHash, p->id(), interval)) + { + JLOG(journal_.trace()) << "Sending " << protocolMessageName(type) << " message to [" + << p->id() << "]: " << messageHash; + p->send(packet); + } + else + JLOG(journal_.debug()) + << "Suppressing sending duplicate " << protocolMessageName(type) + << " message to [" << p->id() << "]: " << messageHash; + } } } diff --git a/src/xrpld/overlay/detail/ProtocolMessage.h b/src/xrpld/overlay/detail/ProtocolMessage.h index 41d42674ba..ef05663ef3 100644 --- a/src/xrpld/overlay/detail/ProtocolMessage.h +++ b/src/xrpld/overlay/detail/ProtocolMessage.h @@ -23,6 +23,12 @@ protocolMessageType(protocol::TMGetLedger const&) return protocol::mtGET_LEDGER; } +inline protocol::MessageType +protocolMessageType(protocol::TMLedgerData const&) +{ + return protocol::mtLEDGER_DATA; +} + inline protocol::MessageType protocolMessageType(protocol::TMReplayDeltaRequest const&) { @@ -432,3 +438,63 @@ invokeProtocolMessage(Buffers const& buffers, Handler& handler, std::size_t& hin } } // namespace xrpl + +namespace protocol { + +template +void +hash_append(Hasher& h, TMGetLedger const& msg) +{ + using beast::hash_append; + using namespace xrpl; + hash_append(h, safe_cast(protocolMessageType(msg))); + hash_append(h, safe_cast(msg.itype())); + if (msg.has_ltype()) + hash_append(h, safe_cast(msg.ltype())); + + if (msg.has_ledgerhash()) + hash_append(h, msg.ledgerhash()); + + if (msg.has_ledgerseq()) + hash_append(h, msg.ledgerseq()); + + for (auto const& nodeId : msg.nodeids()) + hash_append(h, nodeId); + hash_append(h, msg.nodeids_size()); + + // Do NOT include the request cookie. It does not affect the content of the + // request, but only where to route the results. + // if (msg.has_requestcookie()) + // hash_append(h, msg.requestcookie()); + + if (msg.has_querytype()) + hash_append(h, safe_cast(msg.querytype())); + + if (msg.has_querydepth()) + hash_append(h, msg.querydepth()); +} + +template +void +hash_append(Hasher& h, TMLedgerData const& msg) +{ + using beast::hash_append; + using namespace xrpl; + hash_append(h, safe_cast(protocolMessageType(msg))); + hash_append(h, msg.ledgerhash()); + hash_append(h, msg.ledgerseq()); + hash_append(h, safe_cast(msg.type())); + for (auto const& node : msg.nodes()) + { + hash_append(h, node.nodedata()); + if (node.has_nodeid()) + hash_append(h, node.nodeid()); + } + hash_append(h, msg.nodes_size()); + if (msg.has_requestcookie()) + hash_append(h, msg.requestcookie()); + if (msg.has_error()) + hash_append(h, safe_cast(msg.error())); +} + +} // namespace protocol diff --git a/src/xrpld/overlay/detail/ProtocolVersion.cpp b/src/xrpld/overlay/detail/ProtocolVersion.cpp index 1a55030cd4..e4a46ab29f 100644 --- a/src/xrpld/overlay/detail/ProtocolVersion.cpp +++ b/src/xrpld/overlay/detail/ProtocolVersion.cpp @@ -20,6 +20,8 @@ namespace xrpl { constexpr ProtocolVersion const supportedProtocolList[]{ {2, 1}, {2, 2}, + // Adds TMLedgerData::responseCookies and directResponse + {2, 3}, }; // This ugly construct ensures that supportedProtocolList is sorted in strictly