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