From e1615671d02d8ff58f29278fb969a5013d510ba3 Mon Sep 17 00:00:00 2001 From: Brett Boston Date: Fri, 17 Jul 2026 16:46:37 -0700 Subject: [PATCH] Add HAVE_TX_SET overlay message for targeted tx set fetching This PR adds a new message type `HAVE_TX_SET` to allow nodes to communicate that they possess a given transaction set. This improves transaction set fetching with parallel transaction set downloading, as parallel transaction set downloading breaks the assumption that SCP message relayers possess all transaction sets referenced in a given SCP message. --- docs/metrics.md | 5 + src/herder/Herder.h | 6 + src/herder/HerderImpl.cpp | 25 +- src/herder/HerderImpl.h | 2 + src/herder/PendingEnvelopes.cpp | 86 +++- src/herder/PendingEnvelopes.h | 23 + src/overlay/ItemFetcher.cpp | 65 ++- src/overlay/ItemFetcher.h | 21 +- src/overlay/OverlayMetrics.cpp | 14 + src/overlay/OverlayMetrics.h | 7 + src/overlay/Peer.cpp | 46 +- src/overlay/Peer.h | 18 + src/overlay/Tracker.cpp | 137 +++++- src/overlay/Tracker.h | 32 +- src/overlay/test/ItemFetcherTests.cpp | 613 ++++++++++++++++++++++++-- src/overlay/test/TrackerTests.cpp | 10 +- src/protocol-curr/xdr | 2 +- 17 files changed, 1037 insertions(+), 75 deletions(-) diff --git a/docs/metrics.md b/docs/metrics.md index 9c5270b66b..2da971c9b9 100644 --- a/docs/metrics.md +++ b/docs/metrics.md @@ -159,6 +159,11 @@ overlay.inbound.live | counter | number of live inbound c overlay.outbound-queue. | timer | time traffic sits in flow-controlled queues overlay.outbound-queue.drop- | meter | number of messages dropped from flow-controlled queues overlay.item-fetcher.next-peer | meter | ask for item past the first one +overlay.item-fetcher.claim-ask | meter | fetch ask targeted a peer believed to hold the item +overlay.item-fetcher.claim-dropped | meter | HAVE_TX_SET dropped at admission: sending peer exceeded its budget for the current window +overlay.item-fetcher.claim-grace-wait | timer | time a tx set fetch waited from tracker creation to its first ask +overlay.item-fetcher.claim-grace-satisfied | meter | tx set fetch's first ask targeted a believed holder +overlay.item-fetcher.claim-grace-expired | meter | tx set fetch's first ask fell back to an SCP message relayer or random peer overlay.memory.flood-known | counter | number of known flooded entries overlay.message.broadcast | meter | message broadcasted overlay.message.read | meter | message received diff --git a/src/herder/Herder.h b/src/herder/Herder.h index e8b9392982..8d05d18b03 100644 --- a/src/herder/Herder.h +++ b/src/herder/Herder.h @@ -144,6 +144,12 @@ class Herder virtual bool recvSCPQuorumSet(Hash const& hash, SCPQuorumSet const& qset) = 0; virtual bool recvTxSet(Hash const& hash, TxSetXDRFrameConstPtr txset) = 0; + + // A peer announced that it has the tx set with `hash` + virtual void recvHaveTxSet(Hash const& hash, Peer::pointer peer) = 0; + + // Returns true iff the current protocol allows empty-tx-set values + virtual bool protocolAllowsEmptyTxSetValues() const = 0; // We are learning about a new transaction. #ifdef BUILD_TESTS // `isLoadgenTx` is true if the transaction was generated by the load diff --git a/src/herder/HerderImpl.cpp b/src/herder/HerderImpl.cpp index f30e2d7140..185907db2e 100644 --- a/src/herder/HerderImpl.cpp +++ b/src/herder/HerderImpl.cpp @@ -1018,8 +1018,11 @@ HerderImpl::sendSCPStateToPeer(uint32 ledgerSeq, Peer::pointer peer) bool log = true; auto maxSlots = Herder::LEDGER_VALIDITY_BRACKET; - auto sendSlot = [weakPeer = std::weak_ptr(peer)](SCPEnvelope const& e, - bool log) { + // Record of which tx sets we've already sent HAVE_TX_SET messages for to + // `peer` + auto claimedTxSets = std::make_shared>(); + auto sendSlot = [this, claimedTxSets, weakPeer = std::weak_ptr(peer)]( + SCPEnvelope const& e, bool log) { // If in the process of shutting down, exit early auto peerPtr = weakPeer.lock(); if (!peerPtr) @@ -1027,6 +1030,11 @@ HerderImpl::sendSCPStateToPeer(uint32 ledgerSeq, Peer::pointer peer) return false; } + // Tell the peer which of the referenced tx sets we hold. Sent before + // the envelope, so the claim is already buffered when the envelope + // triggers a fetch. + mPendingEnvelopes.sendHaveTxSetClaims(e, peerPtr, *claimedTxSets); + StellarMessage m; m.type(SCP_MESSAGE); m.envelope() = e; @@ -1545,6 +1553,19 @@ HerderImpl::recvTxSet(Hash const& hash, TxSetXDRFrameConstPtr txset) return mPendingEnvelopes.recvTxSet(hash, txset); } +void +HerderImpl::recvHaveTxSet(Hash const& hash, Peer::pointer peer) +{ + ZoneScoped; + mPendingEnvelopes.recvHaveTxSet(hash, peer); +} + +bool +HerderImpl::protocolAllowsEmptyTxSetValues() const +{ + return mHerderSCPDriver.protocolAllowsEmptyTxSetValues(); +} + void HerderImpl::peerDoesntHave(MessageType type, uint256 const& itemID, Peer::pointer peer) diff --git a/src/herder/HerderImpl.h b/src/herder/HerderImpl.h index ceb09b3d22..276fb6e3bc 100644 --- a/src/herder/HerderImpl.h +++ b/src/herder/HerderImpl.h @@ -153,6 +153,8 @@ class HerderImpl : public Herder bool recvSCPQuorumSet(Hash const& hash, SCPQuorumSet const& qset) override; bool recvTxSet(Hash const& hash, TxSetXDRFrameConstPtr txset) override; + void recvHaveTxSet(Hash const& hash, Peer::pointer peer) override; + bool protocolAllowsEmptyTxSetValues() const override; void peerDoesntHave(MessageType type, uint256 const& itemID, Peer::pointer peer) override; TxSetResult getTxSet(Hash const& hash) override; diff --git a/src/herder/PendingEnvelopes.cpp b/src/herder/PendingEnvelopes.cpp index d67d8b91be..ada3473148 100644 --- a/src/herder/PendingEnvelopes.cpp +++ b/src/herder/PendingEnvelopes.cpp @@ -30,11 +30,15 @@ PendingEnvelopes::PendingEnvelopes(Application& app, HerderImpl& herder) , mHerder(herder) , mQsetCache(QSET_CACHE_SIZE) , mTxSetFetcher( - app, [](Peer::pointer peer, Hash hash) { peer->sendGetTxSet(hash); }) - , mQuorumSetFetcher(app, [](Peer::pointer peer, - Hash hash) { peer->sendGetQuorumSet(hash); }) + app, [](Peer::pointer peer, Hash hash) { peer->sendGetTxSet(hash); }, + ItemFetcherKind::TxSet) + , mQuorumSetFetcher( + app, + [](Peer::pointer peer, Hash hash) { peer->sendGetQuorumSet(hash); }, + ItemFetcherKind::QuorumSet) , mTxSetCache(TXSET_CACHE_SIZE) , mValueSizeCache(TXSET_CACHE_SIZE + QSET_CACHE_SIZE) + , mAnnouncedTxSets(TXSET_CACHE_SIZE) , mRebuildQuorum(true) , mQuorumTracker(mApp.getConfig().NODE_SEED.getPublicKey()) , mProcessedCount( @@ -255,6 +259,82 @@ PendingEnvelopes::addTxSet(Hash const& hash, uint64 lastSeenSlotIndex, putTxSet(hash, lastSeenSlotIndex, txset); mTxSetFetcher.recv(hash, mFetchTxSetTimer); + + if (lastSeenSlotIndex != 0) + { + // Announce tx set, so long as it was obtained via a live consensus path + // (not restored from the database; we would have already announced + // those) + maybeAnnounceHaveTxSet(hash); + } +} + +void +PendingEnvelopes::maybeAnnounceHaveTxSet(Hash const& hash) +{ + ZoneScoped; + + if (mAnnouncedTxSets.exists(hash)) + { + return; + } + mAnnouncedTxSets.put(hash, true); + + auto msg = std::make_shared(); + msg->type(HAVE_TX_SET); + msg->haveTxSet().txSetHash = hash; + + for (auto const& peer : mApp.getOverlayManager().getAuthenticatedPeers()) + { + // Send HAVE_TX_SET only if the peer supports it + if (peer.second->getRemoteOverlayVersion() >= + Peer::FIRST_OVERLAY_VERSION_SUPPORTING_HAVE_TX_SET) + { + peer.second->sendMessage(msg); + } + } +} + +void +PendingEnvelopes::sendHaveTxSetClaims(SCPEnvelope const& env, + Peer::pointer const& peer, + UnorderedSet& alreadyClaimed) +{ + ZoneScoped; + + // Peers on pre-HAVE_TX_SET overlay versions cannot parse the message + if (peer->getRemoteOverlayVersion() < + Peer::FIRST_OVERLAY_VERSION_SUPPORTING_HAVE_TX_SET) + { + return; + } + + auto maybeHashes = getTxSetHashes(env); + if (!maybeHashes) + { + return; + } + + for (auto const& hash : *maybeHashes) + { + if (hash == Herder::EMPTY_TX_SET_HASH || !hasTxSet(hash) || + !alreadyClaimed.insert(hash).second) + { + continue; + } + + auto msg = std::make_shared(); + msg->type(HAVE_TX_SET); + msg->haveTxSet().txSetHash = hash; + peer->sendMessage(msg); + } +} + +void +PendingEnvelopes::recvHaveTxSet(Hash const& hash, Peer::pointer peer) +{ + ZoneScoped; + mTxSetFetcher.peerClaimsItem(hash, peer); } bool diff --git a/src/herder/PendingEnvelopes.h b/src/herder/PendingEnvelopes.h index a4bb1ff12a..394402af5c 100644 --- a/src/herder/PendingEnvelopes.h +++ b/src/herder/PendingEnvelopes.h @@ -78,6 +78,9 @@ class PendingEnvelopes // keep track of txset/qset hash -> size pairs for quick access RandomEvictionCache mValueSizeCache; + // tx set hashes already announced to peers via HAVE_TX_SET + RandomEvictionCache mAnnouncedTxSets; + bool mRebuildQuorum; QuorumTracker mQuorumTracker; @@ -118,6 +121,10 @@ class PendingEnvelopes void recordReceivedCost(SCPEnvelope const& env); + // Announce possession of the tx set with `hash` to all authenticated + // peers if not already announced + void maybeAnnounceHaveTxSet(Hash const& hash); + UnorderedMap getCostPerValidator(uint64 slotIndex) const; // stops all pending downloads for slots outside the range @@ -185,6 +192,22 @@ class PendingEnvelopes */ bool recvTxSet(Hash const& hash, TxSetXDRFrameConstPtr txset); + /** + * A peer announced that it has the tx set identified by @p hash. Records + * the claim with the tx set fetcher so an active fetch can target that + * peer. No-op if the tx set is not being fetched. + */ + void recvHaveTxSet(Hash const& hash, Peer::pointer peer); + + /** + * Send HAVE_TX_SET to @p peer for every tx set referenced by @p env that + * is held in memory. + * Adds the hashes of all claims sent to @p peer to @p alreadyClaimed. Does + * not send claims for any hashes already present in @p alreadyClaimed. + */ + void sendHaveTxSetClaims(SCPEnvelope const& env, Peer::pointer const& peer, + UnorderedSet& alreadyClaimed); + // Returns true if the tx set is available locally (either in cache or // is an empty-tx-set hash which doesn't need fetching). bool hasTxSet(Hash const& hash) const; diff --git a/src/overlay/ItemFetcher.cpp b/src/overlay/ItemFetcher.cpp index f2bf867c8c..76b8c9c0c4 100644 --- a/src/overlay/ItemFetcher.cpp +++ b/src/overlay/ItemFetcher.cpp @@ -14,8 +14,16 @@ namespace stellar { -ItemFetcher::ItemFetcher(Application& app, AskPeer askPeer) - : mApp(app), mAskPeer(askPeer) +// Cap on the number of distinct hashes for which we buffer early HAVE_TX_SET +// claims (claims for items not yet being tracked). +static size_t const BUFFERED_CLAIMS_CACHE_SIZE = 1000; + +ItemFetcher::ItemFetcher(Application& app, AskPeer askPeer, + ItemFetcherKind kind) + : mApp(app) + , mAskPeer(askPeer) + , mKind(kind) + , mBufferedClaims(BUFFERED_CLAIMS_CACHE_SIZE) { } @@ -28,10 +36,27 @@ ItemFetcher::fetch(Hash const& itemHash, SCPEnvelope const& envelope) if (entryIt == mTrackers.end()) { // not being tracked TrackerPtr tracker = - std::make_shared(mApp, itemHash, mAskPeer); + std::make_shared(mApp, itemHash, mAskPeer, mKind); mTrackers[itemHash] = tracker; tracker->listen(envelope); + + // Seed any HAVE_TX_SET claims that arrived before this tracker existed + // so the first ask can target a known holder rather than blind-asking. + auto* buffered = mBufferedClaims.maybeGet(itemHash); + if (buffered) + { + for (auto const& [nodeID, weak] : *buffered) + { + if (auto peer = weak.lock()) + { + tracker->seedClaim(peer); + } + } + // Clear the consumed claims + buffered->clear(); + } + tracker->tryNextPeer(); } else @@ -156,6 +181,33 @@ ItemFetcher::doesntHave(Hash const& itemHash, Peer::pointer peer) } } +void +ItemFetcher::peerClaimsItem(Hash const& itemHash, Peer::pointer peer) +{ + ZoneScoped; + auto const& iter = mTrackers.find(itemHash); + if (iter != mTrackers.end()) + { + iter->second->peerClaims(peer); + } + else + { + // Not yet tracking this item. Buffer the claim so the tracker can + // be seeded when it is created + auto* buffered = mBufferedClaims.maybeGet(itemHash); + if (buffered) + { + (*buffered)[peer->getPeerID()] = peer; + } + else + { + mBufferedClaims.put(itemHash, + UnorderedMap>{ + {peer->getPeerID(), peer}}); + } + } +} + void ItemFetcher::recv(Hash itemHash, medida::Timer& timer) { @@ -197,5 +249,12 @@ ItemFetcher::getTracker(Hash const& h) } return it->second; } + +size_t +ItemFetcher::getNumBufferedClaims(Hash const& itemHash) +{ + auto* buffered = mBufferedClaims.maybeGet(itemHash); + return buffered ? buffered->size() : 0; +} #endif } diff --git a/src/overlay/ItemFetcher.h b/src/overlay/ItemFetcher.h index c57fd56994..cb74602ac4 100644 --- a/src/overlay/ItemFetcher.h +++ b/src/overlay/ItemFetcher.h @@ -5,11 +5,15 @@ #pragma once #include "overlay/Peer.h" +#include "overlay/Tracker.h" #include "util/NonCopyable.h" +#include "util/RandomEvictionCache.h" #include "util/Timer.h" +#include "util/UnorderedMap.h" #include #include #include +#include namespace medida { @@ -20,7 +24,6 @@ class Timer; namespace stellar { -class Tracker; class TxSetXDRFrame; struct SCPQuorumSet; using SCPQuorumSetPtr = std::shared_ptr; @@ -41,9 +44,11 @@ class ItemFetcher : private NonMovableOrCopyable using TrackerPtr = std::shared_ptr; /** - * Create ItemFetcher that fetches data using @p askPeer delegate. + * Create ItemFetcher that fetches data using @p askPeer delegate. @p kind + * labels which fetcher this is. */ - explicit ItemFetcher(Application& app, AskPeer askPeer); + explicit ItemFetcher(Application& app, AskPeer askPeer, + ItemFetcherKind kind); /** * Fetch data identified by @p hash and needed by @p envelope. Multiple @@ -92,6 +97,11 @@ class ItemFetcher : private NonMovableOrCopyable */ void doesntHave(Hash const& itemHash, Peer::pointer peer); + /** + * Record that @p peer claims to have the data identified by @p itemHash. + */ + void peerClaimsItem(Hash const& itemHash, Peer::pointer peer); + /** * Called when data with given @p itemHash was received. All envelopes * added before with @see fetch and the same @p itemHash will be resent @@ -101,6 +111,7 @@ class ItemFetcher : private NonMovableOrCopyable #ifdef BUILD_TESTS std::shared_ptr getTracker(Hash const& h); + size_t getNumBufferedClaims(Hash const& itemHash); #endif protected: @@ -113,5 +124,9 @@ class ItemFetcher : private NonMovableOrCopyable private: AskPeer mAskPeer; + ItemFetcherKind mKind; + // HAVE_TX_SET claims received for hashes not yet being tracked. + RandomEvictionCache>> + mBufferedClaims; }; } diff --git a/src/overlay/OverlayMetrics.cpp b/src/overlay/OverlayMetrics.cpp index 85bb88e0e0..108ee5cdc1 100644 --- a/src/overlay/OverlayMetrics.cpp +++ b/src/overlay/OverlayMetrics.cpp @@ -37,6 +37,16 @@ OverlayMetrics::OverlayMetrics(Application& app) , mItemFetcherNextPeer(app.getMetrics().NewMeter( {"overlay", "item-fetcher", "next-peer"}, "item-fetcher")) + , mItemFetcherClaimAsk(app.getMetrics().NewMeter( + {"overlay", "item-fetcher", "claim-ask"}, "item-fetcher")) + , mItemFetcherClaimDropped(app.getMetrics().NewMeter( + {"overlay", "item-fetcher", "claim-dropped"}, "item-fetcher")) + , mItemFetcherClaimGraceWait(app.getMetrics().NewTimer( + {"overlay", "item-fetcher", "claim-grace-wait"})) + , mItemFetcherClaimGraceSatisfied(app.getMetrics().NewMeter( + {"overlay", "item-fetcher", "claim-grace-satisfied"}, "item-fetcher")) + , mItemFetcherClaimGraceExpired(app.getMetrics().NewMeter( + {"overlay", "item-fetcher", "claim-grace-expired"}, "item-fetcher")) , mRecvErrorTimer(app.getMetrics().NewTimer({"overlay", "recv", "error"})) , mRecvHelloTimer(app.getMetrics().NewTimer({"overlay", "recv", "hello"})) @@ -82,6 +92,8 @@ OverlayMetrics::OverlayMetrics(Application& app) app.getMetrics().NewTimer({"overlay", "recv", "flood-advert"})) , mRecvFloodDemandTimer( app.getMetrics().NewTimer({"overlay", "recv", "flood-demand"})) + , mRecvHaveTxSetTimer( + app.getMetrics().NewTimer({"overlay", "recv", "have-tx-set"})) , mRecvTxBatchTimer( app.getMetrics().NewTimer({"overlay", "recv", "tx-batch"})) @@ -144,6 +156,8 @@ OverlayMetrics::OverlayMetrics(Application& app) {"overlay", "send", "flood-advert"}, "message")) , mSendFloodDemandMeter(app.getMetrics().NewMeter( {"overlay", "send", "flood-demand"}, "message")) + , mSendHaveTxSetMeter(app.getMetrics().NewMeter( + {"overlay", "send", "have-tx-set"}, "message")) , mMessagesDemanded(app.getMetrics().NewMeter( {"overlay", "flood", "demanded"}, "message")) , mMessagesFulfilledMeter(app.getMetrics().NewMeter( diff --git a/src/overlay/OverlayMetrics.h b/src/overlay/OverlayMetrics.h index 933031fbe0..6a2494c69b 100644 --- a/src/overlay/OverlayMetrics.h +++ b/src/overlay/OverlayMetrics.h @@ -41,6 +41,11 @@ struct OverlayMetrics medida::Timer& mConnectionFloodThrottle; medida::Meter& mItemFetcherNextPeer; + medida::Meter& mItemFetcherClaimAsk; + medida::Meter& mItemFetcherClaimDropped; + medida::Timer& mItemFetcherClaimGraceWait; + medida::Meter& mItemFetcherClaimGraceSatisfied; + medida::Meter& mItemFetcherClaimGraceExpired; medida::Timer& mRecvErrorTimer; medida::Timer& mRecvHelloTimer; @@ -73,6 +78,7 @@ struct OverlayMetrics medida::Timer& mRecvFloodAdvertTimer; medida::Timer& mRecvFloodDemandTimer; + medida::Timer& mRecvHaveTxSetTimer; medida::Timer& mRecvTxBatchTimer; medida::Timer& mMessageDelayInWriteQueueTimer; @@ -108,6 +114,7 @@ struct OverlayMetrics medida::Meter& mSendFloodAdvertMeter; medida::Meter& mSendFloodDemandMeter; + medida::Meter& mSendHaveTxSetMeter; medida::Meter& mMessagesDemanded; medida::Meter& mMessagesFulfilledMeter; medida::Meter& mBannedMessageUnfulfilledMeter; diff --git a/src/overlay/Peer.cpp b/src/overlay/Peer.cpp index db6a4121a4..250d96f548 100644 --- a/src/overlay/Peer.cpp +++ b/src/overlay/Peer.cpp @@ -379,8 +379,6 @@ Peer::startRecurrentTimer() releaseAssert(threadIsMain()); RECURSIVE_LOCK_GUARD(mStateMutex, guard); - constexpr std::chrono::seconds RECURRENT_TIMER_PERIOD(5); - if (shouldAbort(guard)) { return; @@ -443,6 +441,9 @@ Peer::recurrentTimerExpired(asio::error_code const& error) if (!error) { + // Reset the HAVE_TX_SET budget + mHaveTxSetAdmitted.store(0, std::memory_order_relaxed); + auto now = mAppConnector.now(); auto timeout = getIOTimeout(); auto stragglerTimeout = std::chrono::seconds( @@ -821,6 +822,9 @@ Peer::msgSummary(StellarMessage const& msg) return "FLODADVERT"; case FLOOD_DEMAND: return "FLOODDEMAND"; + case HAVE_TX_SET: + return fmt::format(FMT_STRING("HAVETXSET {}"), + hexAbbrev(msg.haveTxSet().txSetHash)); } return "UNKNOWN"; } @@ -894,6 +898,9 @@ Peer::sendMessage(std::shared_ptr msg, bool log) case FLOOD_DEMAND: mOverlayMetrics.mSendFloodDemandMeter.Mark(); break; + case HAVE_TX_SET: + mOverlayMetrics.mSendHaveTxSetMeter.Mark(); + break; }; releaseAssert(mFlowControl); @@ -1099,6 +1106,18 @@ Peer::recvAuthenticatedMessage(AuthenticatedMessage&& msg) cat = "SCP"; break; + case HAVE_TX_SET: + if (mHaveTxSetAdmitted.load(std::memory_order_relaxed) >= + HAVE_TX_SET_MAX_PER_PERIOD) + { + // Drop HAVE_TX_SET message if over the limit + mOverlayMetrics.mItemFetcherClaimDropped.Mark(); + return true; + } + mHaveTxSetAdmitted.fetch_add(1, std::memory_order_relaxed); + cat = "SCP"; + break; + default: cat = "MISC"; } @@ -1404,6 +1423,13 @@ Peer::recvRawMessage(std::shared_ptr msgTracker) auto t = mOverlayMetrics.mRecvFloodDemandTimer.TimeScope(); recvFloodDemand(stellarMsg); } + break; + + case HAVE_TX_SET: + { + auto t = mOverlayMetrics.mRecvHaveTxSetTimer.TimeScope(); + recvHaveTxSet(stellarMsg); + } } } @@ -2101,6 +2127,22 @@ Peer::recvFloodDemand(StellarMessage const& msg) shared_from_this()); } +void +Peer::recvHaveTxSet(StellarMessage const& msg) +{ + releaseAssert(threadIsMain()); + if (getRemoteOverlayVersion() < + FIRST_OVERLAY_VERSION_SUPPORTING_HAVE_TX_SET) + { + // The peer advertised an overlay version that predates HAVE_TX_SET, + // yet sent one anyway. + sendErrorAndDrop(ERR_MISC, "HAVE_TX_SET from old overlay version"); + return; + } + mAppConnector.getHerder().recvHaveTxSet(msg.haveTxSet().txSetHash, + shared_from_this()); +} + Peer::PeerMetrics::PeerMetrics(VirtualClock::time_point connectedTime) : mMessageRead(0) , mMessageWrite(0) diff --git a/src/overlay/Peer.h b/src/overlay/Peer.h index 3435338fe3..b11851691c 100644 --- a/src/overlay/Peer.h +++ b/src/overlay/Peer.h @@ -98,6 +98,15 @@ class Peer : public std::enable_shared_from_this, uint32_t mNumQueries{0}; }; + // Cadence of the per-peer recurrent timer + static constexpr std::chrono::seconds RECURRENT_TIMER_PERIOD{5}; + + // Max HAVE_TX_SET messages accepted per peer per RECURRENT_TIMER_PERIOD + static constexpr uint32_t HAVE_TX_SET_MAX_PER_PERIOD = 32; + + // First overlay protocol version supporting the HAVE_TX_SET message. + static constexpr uint32_t FIRST_OVERLAY_VERSION_SUPPORTING_HAVE_TX_SET = 42; + static inline int format_as(PeerState const& s) { @@ -275,6 +284,8 @@ class Peer : public std::enable_shared_from_this, QueryInfo mQSetQueryInfo; QueryInfo mTxSetQueryInfo; QueryInfo mSCPStateQueryInfo; + // HAVE_TX_SET messages accepted in the current admission window + std::atomic mHaveTxSetAdmitted{0}; bool mPeersReceived{false}; static Hash pingIDfromTimePoint(VirtualClock::time_point const& tp); @@ -311,6 +322,7 @@ class Peer : public std::enable_shared_from_this, void recvGetSCPState(StellarMessage const& msg); void recvFloodAdvert(StellarMessage const& msg); void recvFloodDemand(StellarMessage const& msg); + void recvHaveTxSet(StellarMessage const& msg); void sendHello(); void sendAuth(); @@ -508,6 +520,12 @@ class Peer : public std::enable_shared_from_this, releaseAssert(threadIsMain()); return mSCPStateQueryInfo.mNumQueries; } + + uint32_t + getHaveTxSetAdmittedForTesting() const + { + return mHaveTxSetAdmitted.load(std::memory_order_relaxed); + } #endif // Public thread-safe methods that access Peer's state diff --git a/src/overlay/Tracker.cpp b/src/overlay/Tracker.cpp index 2ea8acd02b..97f870da34 100644 --- a/src/overlay/Tracker.cpp +++ b/src/overlay/Tracker.cpp @@ -19,20 +19,30 @@ namespace stellar { -std::chrono::milliseconds const Tracker::MS_TO_WAIT_FOR_FETCH_REPLY{1500}; +std::chrono::milliseconds const Tracker::MS_TO_WAIT_FOR_FETCH_PROGRESS{1500}; int const Tracker::MAX_REBUILD_FETCH_LIST = 10; -Tracker::Tracker(Application& app, Hash const& hash, AskPeer& askPeer) +Tracker::Tracker(Application& app, Hash const& hash, AskPeer& askPeer, + ItemFetcherKind kind) : mAskPeer(askPeer) , mApp(app) , mNumListRebuild(0) , mTimer(app) , mItemHash(hash) + , mKind(kind) , mTryNextPeer( app.getOverlayManager().getOverlayMetrics().mItemFetcherNextPeer) , mFetchTime("fetch-" + hexAbbrev(hash), LogSlowExecution::Mode::MANUAL) { releaseAssert(mAskPeer); + + if (mKind == ItemFetcherKind::TxSet && + mApp.getHerder().protocolAllowsEmptyTxSetValues()) + { + mGraceEnabled = true; + mGraceStart = mApp.getClock().now(); + mGraceDeadline = mGraceStart + MS_TO_WAIT_FOR_FETCH_PROGRESS; + } } Tracker::~Tracker() @@ -87,10 +97,31 @@ Tracker::doesntHave(Peer::pointer peer) if (mLastAskedPeer == peer) { CLOG_TRACE(Overlay, "Does not have {}", hexAbbrev(mItemHash)); + // Any claim this peer made was wrong. + mClaimingPeers.erase(peer); + tryNextPeer(); + } +} + +void +Tracker::peerClaims(Peer::pointer peer) +{ + ZoneScoped; + mClaimingPeers.insert(peer); + if (!mLastAskedPeer) + { + // No ask is outstanding. Act on the claim immediately + mTimer.cancel(); tryNextPeer(); } } +void +Tracker::seedClaim(Peer::pointer peer) +{ + mClaimingPeers.insert(peer); +} + void Tracker::tryNextPeer() { @@ -120,9 +151,10 @@ Tracker::tryNextPeer() // // We want to bias the candidates set towards peers that are close to us in // terms of network latency, so we repeatedly lower a "nearness threshold" - // in units of 500ms (1/3 of the MS_TO_WAIT_FOR_FETCH_REPLY) until we have a - // "closest peers" bucket that we have at least one peer for, and keep all - // the peers in that bucket, and then (later) randomly select from it. + // in units of 500ms (1/3 of the MS_TO_WAIT_FOR_FETCH_PROGRESS) until we + // have a "closest peers" bucket that we have at least one peer for, and + // keep all the peers in that bucket, and then (later) randomly select from + // it. // // if the map of peers passed in is for peers that claim to have the data we // need, `peersHave` is also set to true. in this case, the candidate list @@ -138,7 +170,8 @@ Tracker::tryNextPeer() auto& p = mp.second; if (canAskPeer(p, peersHave)) { - int64 GROUPSIZE_MS = (MS_TO_WAIT_FOR_FETCH_REPLY.count() / 3); + int64 GROUPSIZE_MS = + (MS_TO_WAIT_FOR_FETCH_PROGRESS.count() / 3); int64 plat = p->getPing().count() / GROUPSIZE_MS; if (plat < curBest) { @@ -154,7 +187,18 @@ Tracker::tryNextPeer() } }; - // build the set of peers we didn't ask yet that have this envelope + // Whether a peer that relayed an SCP envelope can be assumed to possess + // the item it references. For quorum sets, always: an envelope is never + // relayed without its qset. For tx sets, this function returns true prior + // to the protocol that allows empty-tx-set values, as well as for peers + // running pre-HAVE_TX_SET overlay versions + auto relayerImpliesHave = [&](Peer::pointer const& p) { + return mKind == ItemFetcherKind::QuorumSet || + !mApp.getHerder().protocolAllowsEmptyTxSetValues() || + p->getRemoteOverlayVersion() < + Peer::FIRST_OVERLAY_VERSION_SUPPORTING_HAVE_TX_SET; + }; + std::map newPeersWithEnvelope; for (auto const& e : mWaitingEnvelopes) { @@ -162,20 +206,61 @@ Tracker::tryNextPeer() for (auto pit = s.begin(); pit != s.end(); ++pit) { auto& p = *pit; - if (canAskPeer(p, true)) + if (relayerImpliesHave(p)) + { + mClaimingPeers.insert(p); + } + else if (canAskPeer(p, false)) { newPeersWithEnvelope.emplace(p->getPeerID(), *pit); } } } - bool peerWithEnvelopeSelected = !newPeersWithEnvelope.empty(); - if (peerWithEnvelopeSelected) + // Ask-able peers believed to hold the item + std::map claimingPeers; + for (auto const& p : mClaimingPeers) { - procPeers(newPeersWithEnvelope, true); + if (canAskPeer(p, true)) + { + claimingPeers.emplace(p->getPeerID(), p); + } + } + + // Claim grace period: with no peer known to hold the tx set, prefer waiting + // a bounded time for a HAVE_TX_SET claim over blind-asking a peer that + // likely does not have it yet. An arriving claim preempts the wait. If the + // grace period has expired we fall through to the relayer/random tiers as + // usual. + auto const now = mApp.getClock().now(); + if (claimingPeers.empty() && mGraceEnabled && now < mGraceDeadline) + { + auto const wait = std::chrono::duration_cast( + mGraceDeadline - now); + mTimer.expires_from_now(wait); + mTimer.async_wait([this]() { this->tryNextPeer(); }, + VirtualTimer::onFailureNoop); + return; + } + + bool claimTierSelected = false; + bool selectedPeersHave = false; + if (!claimingPeers.empty()) + { + // Prefer peers who claim to have the item being fetched + claimTierSelected = true; + selectedPeersHave = true; + procPeers(claimingPeers, true); + } + else if (!newPeersWithEnvelope.empty()) + { + // If no peers claim to have the item, fall back on those who claim to + // have at least heard of it via SCP messages + procPeers(newPeersWithEnvelope, false); } else { + // If all else fails, ask a random peer auto& inPeers = mApp.getOverlayManager().getInboundAuthenticatedPeers(); auto& outPeers = mApp.getOverlayManager().getOutboundAuthenticatedPeers(); @@ -200,16 +285,39 @@ Tracker::tryNextPeer() CLOG_TRACE(Overlay, "tryNextPeer {} restarting fetch #{}", hexAbbrev(mItemHash), mNumListRebuild); - nextTry = MS_TO_WAIT_FOR_FETCH_REPLY * + nextTry = MS_TO_WAIT_FOR_FETCH_PROGRESS * std::min(MAX_REBUILD_FETCH_LIST, mNumListRebuild); } else { - mPeersAsked[mLastAskedPeer] = peerWithEnvelopeSelected; + auto& om = mApp.getOverlayManager().getOverlayMetrics(); + if (claimTierSelected) + { + // Record that we requested from a peer who claims to have the data + // we want + om.mItemFetcherClaimAsk.Mark(); + } + // Record the grace period outcome once, at the first ask: did waiting + // for a claim land us on a claimer, or did we fall back to a blind + // ask? mGraceEnabled is only set for TxSet fetches. + if (mGraceEnabled && !mGraceResolved) + { + mGraceResolved = true; + om.mItemFetcherClaimGraceWait.Update(now - mGraceStart); + if (claimTierSelected) + { + om.mItemFetcherClaimGraceSatisfied.Mark(); + } + else + { + om.mItemFetcherClaimGraceExpired.Mark(); + } + } + mPeersAsked[mLastAskedPeer] = selectedPeersHave; CLOG_TRACE(Overlay, "Asking for {} to {}", hexAbbrev(mItemHash), mLastAskedPeer->toString()); mAskPeer(mLastAskedPeer, mItemHash); - nextTry = MS_TO_WAIT_FOR_FETCH_REPLY; + nextTry = MS_TO_WAIT_FOR_FETCH_PROGRESS; } mTimer.expires_from_now(nextTry); @@ -267,6 +375,7 @@ Tracker::cancel() { mTimer.cancel(); mLastSeenSlotIndex = 0; + mClaimingPeers.clear(); } std::chrono::milliseconds diff --git a/src/overlay/Tracker.h b/src/overlay/Tracker.h index 7af1d104f3..7e7cc2e0df 100644 --- a/src/overlay/Tracker.h +++ b/src/overlay/Tracker.h @@ -39,6 +39,13 @@ class Application; using AskPeer = std::function; +// Type of the object the Tracker is fetching +enum class ItemFetcherKind +{ + TxSet, + QuorumSet +}; + class Tracker { private: @@ -49,21 +56,32 @@ class Tracker // keep track of which peer we asked, and if we thought if it had the data // or not at the time std::map mPeersAsked; + // Peers claiming to have the data + UnorderedSet mClaimingPeers; VirtualTimer mTimer; + // Claim-grace bookkeeping (TxSet kind only). While within the grace window + // and with no claiming peer, the fetch waits for a HAVE_TX_SET claim rather + // than blind-asking an SCP message relayer + VirtualClock::time_point mGraceStart; + VirtualClock::time_point mGraceDeadline; + bool mGraceEnabled{false}; + bool mGraceResolved{false}; std::vector> mWaitingEnvelopes; Hash mItemHash; + ItemFetcherKind mKind; medida::Meter& mTryNextPeer; uint64 mLastSeenSlotIndex{0}; LogSlowExecution mFetchTime; public: - static std::chrono::milliseconds const MS_TO_WAIT_FOR_FETCH_REPLY; + static std::chrono::milliseconds const MS_TO_WAIT_FOR_FETCH_PROGRESS; static int const MAX_REBUILD_FETCH_LIST; /** * Create Tracker that tracks data identified by @p hash. @p askPeer * delegate is used to fetch the data. */ - explicit Tracker(Application& app, Hash const& hash, AskPeer& askPeer); + explicit Tracker(Application& app, Hash const& hash, AskPeer& askPeer, + ItemFetcherKind kind); virtual ~Tracker(); /** @@ -137,6 +155,16 @@ class Tracker */ void doesntHave(Peer::pointer peer); + /** + * Peer @p peer claims to have the item this Tracker is trying to fetch + */ + void peerClaims(Peer::pointer peer); + + /** + * Record that @p peer is a claiming peer without triggering an ask. + */ + void seedClaim(Peer::pointer peer); + /** * Called either when @see doesntHave(Peer::pointer) was received or * request to peer timed out. diff --git a/src/overlay/test/ItemFetcherTests.cpp b/src/overlay/test/ItemFetcherTests.cpp index 908f6e27c4..8189e85879 100644 --- a/src/overlay/test/ItemFetcherTests.cpp +++ b/src/overlay/test/ItemFetcherTests.cpp @@ -7,6 +7,7 @@ #include "crypto/SHA.h" #include "herder/Herder.h" #include "herder/HerderImpl.h" +#include "herder/TxSetFrame.h" #include "main/ApplicationImpl.h" #include "overlay/ItemFetcher.h" #include "overlay/OverlayManager.h" @@ -17,7 +18,11 @@ #include "test/TestUtils.h" #include "test/test.h" #include "util/MetricsRegistry.h" +#include "util/ProtocolVersion.h" #include "xdr/Stellar-types.h" +#include +#include +#include namespace stellar { @@ -75,6 +80,87 @@ makeEnvelope(int id) result.statement.pledges.confirm().nPrepared = id; return result; } + +// Shared setup for the fetcher/claim simulation tests +struct FetchTestSim +{ + Simulation::pointer sim; + // Main's app, and its Peer objects + Application::pointer app; + std::vector> peers; + // The remote ends, used to send TO Main, and the remote apps + std::vector> senders; + std::vector remoteApps; + + // `legacyRemote` pins that remote's overlay version below + // FIRST_OVERLAY_VERSION_SUPPORTING_HAVE_TX_SET. + FetchTestSim(size_t numRemotes, + std::optional legacyRemote = std::nullopt) + : sim(std::make_shared( + Simulation::OVER_LOOPBACK, + sha256(getTestConfig().NETWORK_PASSPHRASE))) + , mMainKey(SecretKey::fromSeed(sha256("NODE_SEED_MAIN"))) + { + auto cfgMain = getTestConfig(1); + quieten(cfgMain); + sim->addNode(mMainKey, cfgMain.QUORUM_SET, &cfgMain); + sim->startAllNodes(); + app = sim->getNode(mMainKey.getPublicKey()); + for (size_t i = 0; i < numRemotes; ++i) + { + addRemote(legacyRemote == i); + } + } + + // Add another remote node connected to Main, crank until authenticated, + // and return Main's Peer object for it. + std::shared_ptr + addRemote(bool legacy = false) + { + size_t const i = peers.size(); + auto cfg = getTestConfig(static_cast(2 + i)); + if (legacy) + { + cfg.OVERLAY_PROTOCOL_VERSION = + Peer::FIRST_OVERLAY_VERSION_SUPPORTING_HAVE_TX_SET - 1; + } + quieten(cfg); + auto const key = SecretKey::fromSeed( + sha256("NODE_SEED_REMOTE_" + std::to_string(i))); + sim->addNode(key, cfg.QUORUM_SET, &cfg); + sim->addPendingConnection(mMainKey.getPublicKey(), key.getPublicKey()); + sim->startAllNodes(); + auto conn = sim->getLoopbackConnection(mMainKey.getPublicKey(), + key.getPublicKey()); + releaseAssert(conn); + peers.emplace_back(conn->getInitiator()); + senders.emplace_back(conn->getAcceptor()); + remoteApps.emplace_back(sim->getNode(key.getPublicKey())); + sim->crankUntil( + [&]() { + return std::all_of(peers.begin(), peers.end(), + [](auto const& p) { + return p->isAuthenticatedForTesting(); + }) && + std::all_of(senders.begin(), senders.end(), + [](auto const& p) { + return p->isAuthenticatedForTesting(); + }); + }, + std::chrono::seconds{3}, false); + return peers.back(); + } + + private: + static void + quieten(Config& cfg) + { + cfg.NODE_IS_VALIDATOR = false; + cfg.FORCE_SCP = false; + } + + SecretKey const mMainKey; +}; } TEST_CASE("ItemFetcher fetches", "[overlay][ItemFetcher]") @@ -86,11 +172,14 @@ TEST_CASE("ItemFetcher fetches", "[overlay][ItemFetcher]") std::vector asked; std::vector askedTP; std::vector received; - ItemFetcher itemFetcher(*app, [&](Peer::pointer peer, Hash hash) { - asked.emplace_back(peer); - askedTP.emplace_back(clock.now()); - peer->sendGetQuorumSet(hash); - }); + ItemFetcher itemFetcher( + *app, + [&](Peer::pointer peer, Hash hash) { + asked.emplace_back(peer); + askedTP.emplace_back(clock.now()); + peer->sendGetQuorumSet(hash); + }, + ItemFetcherKind::QuorumSet); auto checkFetchingFor = [&itemFetcher](Hash hash, std::vector envelopes) { @@ -270,7 +359,7 @@ TEST_CASE("ItemFetcher fetches", "[overlay][ItemFetcher]") } }; - crankFor(Tracker::MS_TO_WAIT_FOR_FETCH_REPLY * 2); + crankFor(Tracker::MS_TO_WAIT_FOR_FETCH_PROGRESS * 2); REQUIRE(asked.size() == 2); @@ -341,14 +430,14 @@ TEST_CASE("ItemFetcher fetches", "[overlay][ItemFetcher]") else { REQUIRE(delta >= - Tracker::MS_TO_WAIT_FOR_FETCH_REPLY); + Tracker::MS_TO_WAIT_FOR_FETCH_PROGRESS); } if (i > 0) { auto deltaGroup = refTP - prevGroup; // gap between groups depend on number of retries auto nextTry = - Tracker::MS_TO_WAIT_FOR_FETCH_REPLY * + Tracker::MS_TO_WAIT_FOR_FETCH_PROGRESS * std::min(Tracker::MAX_REBUILD_FETCH_LIST, (static_cast(i - 1)) / 2); REQUIRE(deltaGroup >= nextTry); @@ -382,32 +471,14 @@ TEST_CASE("ItemFetcher fetches", "[overlay][ItemFetcher]") TEST_CASE("next peer strategy", "[overlay][ItemFetcher]") { - auto networkID = sha256(getTestConfig().NETWORK_PASSPHRASE); - auto sim = - std::make_shared(Simulation::OVER_LOOPBACK, networkID); - - auto cfgMain = getTestConfig(1); - auto cfg1 = getTestConfig(2); - auto cfg2 = getTestConfig(3); - - SIMULATION_CREATE_NODE(Main); - SIMULATION_CREATE_NODE(Node1); - SIMULATION_CREATE_NODE(Node2); - sim->addNode(vMainSecretKey, cfgMain.QUORUM_SET, &cfgMain); - - sim->addNode(vNode1SecretKey, cfg1.QUORUM_SET, &cfg1); - sim->addPendingConnection(vMainNodeID, vNode1NodeID); - sim->startAllNodes(); - auto conn1 = sim->getLoopbackConnection(vMainNodeID, vNode1NodeID); - auto peer1 = conn1->getInitiator(); - - auto app = sim->getNode(vMainNodeID); + FetchTestSim s(1); + auto app = s.app; + auto peer1 = s.peers[0]; int askCount = 0; - ItemFetcher itemFetcher(*app, [&](Peer::pointer, Hash) { askCount++; }); - - sim->crankUntil([&]() { return peer1->isAuthenticatedForTesting(); }, - std::chrono::seconds{3}, false); + ItemFetcher itemFetcher( + *app, [&](Peer::pointer, Hash) { askCount++; }, + ItemFetcherKind::QuorumSet); // this causes to fetch from `peer1` as it's the only one // connected @@ -431,14 +502,7 @@ TEST_CASE("next peer strategy", "[overlay][ItemFetcher]") } SECTION("with more peers") { - sim->addNode(vNode2SecretKey, cfg2.QUORUM_SET, &cfg2); - sim->addPendingConnection(vMainNodeID, vNode2NodeID); - sim->startAllNodes(); - auto conn2 = sim->getLoopbackConnection(vMainNodeID, vNode2NodeID); - auto peer2 = conn2->getInitiator(); - - sim->crankUntil([&]() { return peer2->isAuthenticatedForTesting(); }, - std::chrono::seconds{3}, false); + auto peer2 = s.addRemote(); // still connected REQUIRE(peer1->isAuthenticatedForTesting()); @@ -484,4 +548,473 @@ TEST_CASE("next peer strategy", "[overlay][ItemFetcher]") } } } + +TEST_CASE("ItemFetcher claims", "[overlay][ItemFetcher]") +{ + FetchTestSim s(2); + auto sim = s.sim; + auto app = s.app; + auto peer1 = s.peers[0]; + auto peer2 = s.peers[1]; + + int askCount = 0; + ItemFetcher itemFetcher( + *app, [&](Peer::pointer, Hash) { askCount++; }, ItemFetcherKind::TxSet); + + auto& claimAsk = app->getMetrics().NewMeter( + {"overlay", "item-fetcher", "claim-ask"}, "item-fetcher"); + + auto env = makeEnvelope(200); + auto hash = sha256(ByteSlice("200")); + + SECTION("claim with no live tracker is a no-op") + { + itemFetcher.peerClaimsItem(hash, peer1); + REQUIRE(askCount == 0); + } + + SECTION("claim re-enables a missed peer and targets it") + { + itemFetcher.fetch(hash, env); + // Ride out the claim grace period: with no claim known, the first + // (blind) ask fires once the grace expires. + sim->crankUntil([&]() { return askCount == 1; }, + Tracker::MS_TO_WAIT_FOR_FETCH_PROGRESS * 2, false); + auto tracker = itemFetcher.getTracker(hash); + REQUIRE(tracker); + auto first = tracker->getLastAskedPeer(); + REQUIRE(first); + auto second = (first == peer1) ? peer2 : peer1; + + // First peer misses; the fetcher moves on to the second + tracker->doesntHave(first); + REQUIRE(askCount == 2); + REQUIRE(tracker->getLastAskedPeer() == second); + + // A claim from the missed peer while an ask is outstanding is + // recorded but does not interrupt the outstanding ask + tracker->peerClaims(first); + REQUIRE(askCount == 2); + + // When the second peer also misses, the claiming peer is re-asked + // even though it was asked (and missed) before + auto const claimAsksBefore = claimAsk.count(); + tracker->doesntHave(second); + REQUIRE(askCount == 3); + REQUIRE(tracker->getLastAskedPeer() == first); + REQUIRE(claimAsk.count() == claimAsksBefore + 1); + } + + SECTION("claim while idling acts immediately") + { + itemFetcher.fetch(hash, env); + // Ride out the claim grace: with no claim known, the first (blind) + // ask fires once the grace expires. + sim->crankUntil([&]() { return askCount == 1; }, + Tracker::MS_TO_WAIT_FOR_FETCH_PROGRESS * 2, false); + auto tracker = itemFetcher.getTracker(hash); + REQUIRE(tracker); + auto first = tracker->getLastAskedPeer(); + auto second = (first == peer1) ? peer2 : peer1; + + // Both peers miss; with no askable candidates the tracker idles + // waiting for a rebuild + tracker->doesntHave(first); + REQUIRE(askCount == 2); + tracker->doesntHave(second); + REQUIRE(askCount == 2); + REQUIRE(!tracker->getLastAskedPeer()); + + // A claim is acted on immediately, without waiting out the timer + tracker->peerClaims(second); + REQUIRE(askCount == 3); + REQUIRE(tracker->getLastAskedPeer() == second); + } +} + +TEST_CASE("ItemFetcher claim buffer and grace period", "[overlay][ItemFetcher]") +{ + FetchTestSim s(2); + auto app = s.app; + auto peer1 = s.peers[0]; + auto peer2 = s.peers[1]; + + int askCount = 0; + Peer::pointer lastAsked; + ItemFetcher itemFetcher( + *app, + [&](Peer::pointer p, Hash) { + ++askCount; + lastAsked = p; + }, + ItemFetcherKind::TxSet); + + auto& claimAsk = app->getMetrics().NewMeter( + {"overlay", "item-fetcher", "claim-ask"}, "item-fetcher"); +#ifdef CAP_0083 + auto& graceSatisfied = app->getMetrics().NewMeter( + {"overlay", "item-fetcher", "claim-grace-satisfied"}, "item-fetcher"); + auto& graceExpired = app->getMetrics().NewMeter( + {"overlay", "item-fetcher", "claim-grace-expired"}, "item-fetcher"); +#endif + + auto env = makeEnvelope(300); + auto hash = sha256(ByteSlice("300")); + + SECTION("buffered claim seeds the tracker's first ask") + { + // Claim arrives before we are tracking the item: it must be buffered, + // not dropped. + itemFetcher.peerClaimsItem(hash, peer1); + REQUIRE(askCount == 0); + + // When the tracker is created, the buffered claim seeds the claims + // tier so the very first ask targets the claimer. + itemFetcher.fetch(hash, env); + REQUIRE(askCount == 1); + REQUIRE(lastAsked == peer1); + REQUIRE(claimAsk.count() == 1); + } + + SECTION("repeated claims from one peer are deduplicated") + { + // Claims are keyed by peer: re-claiming the same hash cannot grow + // the buffered entry, so a peer's footprint per hash is one slot. + itemFetcher.peerClaimsItem(hash, peer1); + itemFetcher.peerClaimsItem(hash, peer1); + itemFetcher.peerClaimsItem(hash, peer1); + REQUIRE(itemFetcher.getNumBufferedClaims(hash) == 1); + + // A distinct claimer occupies its own slot. + itemFetcher.peerClaimsItem(hash, peer2); + REQUIRE(itemFetcher.getNumBufferedClaims(hash) == 2); + REQUIRE(askCount == 0); + + // Creating the tracker consumes the buffered claims; the first ask + // targets one of the claimers. + itemFetcher.fetch(hash, env); + REQUIRE(askCount == 1); + REQUIRE((lastAsked == peer1 || lastAsked == peer2)); + REQUIRE(claimAsk.count() == 1); + REQUIRE(itemFetcher.getNumBufferedClaims(hash) == 0); + } + +#ifdef CAP_0083 + SECTION("grace period defers the first ask until a claim arrives") + { + itemFetcher.fetch(hash, env); + // With the grace period active and no claimer known, the first ask is + // deferred rather than blind-asking a peer. + REQUIRE(askCount == 0); + + // A claim during the grace window is acted on immediately. + itemFetcher.peerClaimsItem(hash, peer2); + REQUIRE(askCount == 1); + REQUIRE(lastAsked == peer2); + REQUIRE(claimAsk.count() == 1); + REQUIRE(graceSatisfied.count() == 1); + REQUIRE(graceExpired.count() == 0); + } + + SECTION("grace expires and falls back to a blind ask") + { + auto const start = app->getClock().now(); + itemFetcher.fetch(hash, env); + REQUIRE(askCount == 0); + + // No claim arrives; once the grace expires we fall back to asking a + // peer anyway (liveness is preserved). + s.sim->crankUntil([&]() { return askCount == 1; }, + Tracker::MS_TO_WAIT_FOR_FETCH_PROGRESS + + std::chrono::seconds{1}, + false); + auto const elapsed = app->getClock().now() - start; + + REQUIRE(elapsed >= Tracker::MS_TO_WAIT_FOR_FETCH_PROGRESS); + REQUIRE(graceExpired.count() == 1); + REQUIRE(graceSatisfied.count() == 0); + } +#endif // CAP_0083 +} + +TEST_CASE("HAVE_TX_SET admission cap", "[overlay][ItemFetcher]") +{ + FetchTestSim s(2); + auto sim = s.sim; + auto app = s.app; + // Main's view of each remote peer (where the admission counter lives) + auto peer1 = s.peers[0]; + auto peer2 = s.peers[1]; + // The remote ends, used to send claims TO Main + auto sender1 = s.senders[0]; + auto sender2 = s.senders[1]; + + auto& dropped = app->getMetrics().NewMeter( + {"overlay", "item-fetcher", "claim-dropped"}, "item-fetcher"); + + auto claimMsg = [](std::string const& seed) { + StellarMessage msg; + msg.type(HAVE_TX_SET); + msg.haveTxSet().txSetHash = sha256(ByteSlice(seed)); + return std::make_shared(msg); + }; + + uint32_t const Q = Peer::HAVE_TX_SET_MAX_PER_PERIOD; + uint32_t const K = 8; + + // Fill the budget and then some: exactly Q admissions, K drops. + for (uint32_t i = 0; i < Q + K; i++) + { + sender1->sendMessage(claimMsg("cap-" + std::to_string(i)), false); + } + sim->crankUntil([&]() { return dropped.count() >= K; }, + std::chrono::seconds{5}, false); + REQUIRE(dropped.count() == K); + REQUIRE(peer1->getHaveTxSetAdmittedForTesting() == Q); + + // The budget is per peer: another peer's claims are unaffected. + sender2->sendMessage(claimMsg("iso"), false); + sim->crankUntil( + [&]() { return peer2->getHaveTxSetAdmittedForTesting() == 1; }, + std::chrono::seconds{5}, false); + REQUIRE(peer2->getHaveTxSetAdmittedForTesting() == 1); + REQUIRE(dropped.count() == K); + + // Each peer's recurrent timer resets its budget once per period. + sim->crankUntil( + [&]() { + return peer1->getHaveTxSetAdmittedForTesting() == 0 && + peer2->getHaveTxSetAdmittedForTesting() == 0; + }, + Peer::RECURRENT_TIMER_PERIOD * 2, false); + + // Post-reset claims are admitted again. + sender1->sendMessage(claimMsg("post-reset"), false); + sim->crankUntil( + [&]() { return peer1->getHaveTxSetAdmittedForTesting() == 1; }, + std::chrono::seconds{5}, false); + REQUIRE(peer1->getHaveTxSetAdmittedForTesting() == 1); + REQUIRE(dropped.count() == K); +} + +TEST_CASE("HAVE_TX_SET announce respects peer capability", + "[overlay][ItemFetcher]") +{ + // The second remote predates HAVE_TX_SET: it must never be sent the + // message + FetchTestSim s(2, /*legacyRemote=*/1); + auto sim = s.sim; + auto app = s.app; + auto peer1 = s.peers[0]; + auto peer2 = s.peers[1]; + + auto& sent = app->getMetrics().NewMeter({"overlay", "send", "have-tx-set"}, + "message"); + + auto& pe = static_cast(app->getHerder()).getPendingEnvelopes(); + auto txset = TxSetXDRFrame::makeEmpty( + app->getLedgerManager().getLastClosedLedgerHeader()); + auto hash = sha256(ByteSlice("announce")); + + // Obtaining a tx set announces it — but only to the capable peer. + pe.addTxSet(hash, 10, txset); + REQUIRE(sent.count() == 1); + + // Announce-once per hash: re-adding does not re-announce. + pe.addTxSet(hash, 10, txset); + REQUIRE(sent.count() == 1); + + // Tx sets restored from disk at startup (slot 0 sentinel) are not + // announced. + pe.addTxSet(sha256(ByteSlice("restored")), 0, txset); + REQUIRE(sent.count() == 1); +} + +TEST_CASE("HAVE_TX_SET claims accompany SCP state", "[overlay][ItemFetcher]") +{ + // The second remote predates HAVE_TX_SET: it must never be sent claims + FetchTestSim s(2, /*legacyRemote=*/1); + auto sim = s.sim; + auto app = s.app; + auto peer1 = s.peers[0]; + auto peer2 = s.peers[1]; + + auto& sent = app->getMetrics().NewMeter({"overlay", "send", "have-tx-set"}, + "message"); + auto const sentBase = sent.count(); + + auto& pe = static_cast(app->getHerder()).getPendingEnvelopes(); + auto txset = TxSetXDRFrame::makeEmpty( + app->getLedgerManager().getLastClosedLedgerHeader()); + auto heldHash = sha256(ByteSlice("held")); + auto unheldHash = sha256(ByteSlice("unheld")); + // Take possession directly, without addTxSet's announce side effect. + pe.putTxSet(heldHash, 2, txset); + + // An envelope whose nomination votes reference a held set, an unheld + // set, and the empty-tx-set value. + SCPEnvelope env; + env.statement.slotIndex = 2; + env.statement.pledges.type(SCP_ST_NOMINATE); + auto addVote = [&](Hash const& h) { + StellarValue sv; + sv.txSetHash = h; + sv.closeTime = 1; + env.statement.pledges.nominate().votes.emplace_back( + xdr::xdr_to_opaque(sv)); + }; + addVote(heldHash); + addVote(unheldHash); + addVote(Herder::EMPTY_TX_SET_HASH); + + // Only the held set is claimed: the unheld set and the emtpy-tx-set are + // skipped. + UnorderedSet claimed; + pe.sendHaveTxSetClaims(env, peer1, claimed); + REQUIRE(sent.count() == sentBase + 1); + REQUIRE(claimed == UnorderedSet{heldHash}); + + // Within one exchange, further envelopes referencing the same set do not + // re-claim it. + pe.sendHaveTxSetClaims(env, peer1, claimed); + REQUIRE(sent.count() == sentBase + 1); + + // Legacy peers are never sent claims. + UnorderedSet claimedLegacy; + pe.sendHaveTxSetClaims(env, peer2, claimedLegacy); + REQUIRE(sent.count() == sentBase + 1); + REQUIRE(claimedLegacy.empty()); + + // The claims actually arrive at the capable peer, and both connections + // stay up. + auto node1 = s.remoteApps[0]; + sim->crankUntil( + [&]() { + return node1->getMetrics() + .NewTimer({"overlay", "recv", "have-tx-set"}) + .count() >= 1; + }, + std::chrono::seconds{5}, false); + REQUIRE(peer1->isAuthenticatedForTesting()); + REQUIRE(peer2->isAuthenticatedForTesting()); +} + +TEST_CASE("HAVE_TX_SET from unsupporting peer is rejected", + "[overlay][ItemFetcher]") +{ + // The remote advertises a pre-HAVE_TX_SET overlay version. + FetchTestSim s(1, /*legacyRemote=*/0); + auto sim = s.sim; + auto peer1 = s.peers[0]; + auto sender1 = s.senders[0]; + + // The legacy remote sends a HAVE_TX_SET message + StellarMessage msg; + msg.type(HAVE_TX_SET); + msg.haveTxSet().txSetHash = sha256(ByteSlice("violation")); + sender1->sendMessage(std::make_shared(msg), false); + + // Legacy remote is dropped + sim->crankUntil([&]() { return !peer1->isAuthenticatedForTesting(); }, + std::chrono::seconds{5}, false); + REQUIRE(!peer1->isAuthenticatedForTesting()); +} + +#ifdef CAP_0083 +TEST_CASE("relayer possession semantics", "[overlay][ItemFetcher]") +{ + // Once the protocol allows empty-tx-set values, an SCP message relayer + // merely knows OF an item and must claim possession explicitly; a legacy + // relayer keeps the historical implies-possession semantics. + bool legacyPeer = GENERATE(false, true); + + FetchTestSim s(1, legacyPeer ? std::optional{0} : std::nullopt); + auto sim = s.sim; + auto app = s.app; + auto peer1 = s.peers[0]; + + int askCount = 0; + ItemFetcher itemFetcher( + *app, [&](Peer::pointer, Hash) { askCount++; }, ItemFetcherKind::TxSet); + + // With no relayer knowledge, the first ask goes out via the random tier + // (asked without assumed possession). + auto env = makeEnvelope(500); + auto hash = sha256(ByteSlice("500")); + itemFetcher.fetch(hash, env); + // Ride out the claim grace period: with no claim known, the first (blind) + // ask fires once the grace period expires. + sim->crankForAtLeast(Tracker::MS_TO_WAIT_FOR_FETCH_PROGRESS + + std::chrono::milliseconds(100), + false); + auto tracker = itemFetcher.getTracker(hash); + REQUIRE(tracker); + REQUIRE(askCount == 1); + REQUIRE(tracker->getLastAskedPeer() == peer1); + + // Now peer1 relays the envelope. + StellarMessage msg(SCP_MESSAGE); + msg.envelope() = env; + app->getOverlayManager().recvFloodedMsgID(peer1, xdrBlake2(msg)); + + tracker->tryNextPeer(); + if (legacyPeer) + { + // Historical semantics: knowing of the envelope implies having the + // tx set, so the peer is re-asked. + REQUIRE(askCount == 2); + REQUIRE(tracker->getLastAskedPeer() == peer1); + } + else + { + // Knowing OF the item does not re-enable the peer; the fetch idles + // in a rebuild instead. Possession must be claimed explicitly (via + // HAVE_TX_SET), which the claim tests cover. + REQUIRE(askCount == 1); + REQUIRE(!tracker->getLastAskedPeer()); + } +} + +TEST_CASE("legacy relayer bypasses the claim grace period", + "[overlay][ItemFetcher]") +{ + // A relayer whose overlay version predates HAVE_TX_SET provably fetched + // the tx set before relaying, but can never claim possession. Its relay + // is treated as a claim: the fetch targets it immediately instead of + // waiting out the claim grace. + FetchTestSim s(1, /*legacyRemote=*/0); + auto app = s.app; + auto peer1 = s.peers[0]; + + int askCount = 0; + Peer::pointer lastAsked; + ItemFetcher itemFetcher( + *app, + [&](Peer::pointer p, Hash) { + ++askCount; + lastAsked = p; + }, + ItemFetcherKind::TxSet); + + auto& claimAsk = app->getMetrics().NewMeter( + {"overlay", "item-fetcher", "claim-ask"}, "item-fetcher"); + auto& graceSatisfied = app->getMetrics().NewMeter( + {"overlay", "item-fetcher", "claim-grace-satisfied"}, "item-fetcher"); + + // peer1 relays the envelope before the fetch starts. + auto env = makeEnvelope(600); + auto hash = sha256(ByteSlice("600")); + StellarMessage msg(SCP_MESSAGE); + msg.envelope() = env; + app->getOverlayManager().recvFloodedMsgID(peer1, xdrBlake2(msg)); + + // The first ask fires immediately (no grace period delay) and targets the + // legacy relayer. + itemFetcher.fetch(hash, env); + REQUIRE(askCount == 1); + REQUIRE(lastAsked == peer1); + REQUIRE(claimAsk.count() == 1); + REQUIRE(graceSatisfied.count() == 1); +} +#endif // CAP_0083 } diff --git a/src/overlay/test/TrackerTests.cpp b/src/overlay/test/TrackerTests.cpp index 42c912ec81..466ef7ab18 100644 --- a/src/overlay/test/TrackerTests.cpp +++ b/src/overlay/test/TrackerTests.cpp @@ -38,7 +38,7 @@ TEST_CASE("Tracker works", "[overlay][Tracker]") SECTION("empty tracker") { - Tracker t{*app, hash, nullAskPeer}; + Tracker t{*app, hash, nullAskPeer, ItemFetcherKind::TxSet}; REQUIRE(t.size() == 0); REQUIRE(t.empty()); REQUIRE(t.getLastSeenSlotIndex() == 0); @@ -46,7 +46,7 @@ TEST_CASE("Tracker works", "[overlay][Tracker]") SECTION("can listen on envelope") { - Tracker t{*app, hash, nullAskPeer}; + Tracker t{*app, hash, nullAskPeer, ItemFetcherKind::TxSet}; auto env1 = makeEnvelope(1); t.listen(env1); @@ -63,7 +63,7 @@ TEST_CASE("Tracker works", "[overlay][Tracker]") SECTION("listen twice on the same envelope") { - Tracker t{*app, hash, nullAskPeer}; + Tracker t{*app, hash, nullAskPeer, ItemFetcherKind::TxSet}; auto env1 = makeEnvelope(1); t.listen(env1); // this should no-op (idempotent) @@ -76,7 +76,7 @@ TEST_CASE("Tracker works", "[overlay][Tracker]") SECTION("can listen on different envelopes") { - Tracker t{*app, hash, nullAskPeer}; + Tracker t{*app, hash, nullAskPeer, ItemFetcherKind::TxSet}; auto env1 = makeEnvelope(1); auto env2 = makeEnvelope(2); t.listen(env1); @@ -90,7 +90,7 @@ TEST_CASE("Tracker works", "[overlay][Tracker]") SECTION("properly removes old envelopes") { - Tracker t{*app, hash, nullAskPeer}; + Tracker t{*app, hash, nullAskPeer, ItemFetcherKind::TxSet}; auto env1 = makeEnvelope(1); auto env2 = makeEnvelope(2); auto env3 = makeEnvelope(3); diff --git a/src/protocol-curr/xdr b/src/protocol-curr/xdr index 8d1c926a5e..f874506b35 160000 --- a/src/protocol-curr/xdr +++ b/src/protocol-curr/xdr @@ -1 +1 @@ -Subproject commit 8d1c926a5e5d9167d0fc7327b5dad06dc1ec4fd7 +Subproject commit f874506b357421eab455e84640fd1683e9e9d330