Compare commits

...

2 Commits

Author SHA1 Message Date
Bart
0514203d03 refactor: Open the acquisition seams its tests need
TransactionAcquire and InboundLedger are no longer final and take a
retry interval defaulted to the constant each used before, so a test
subclass can drive a whole timeout chain in a fraction of a second.
TimeoutCounter's constructor already asserts an interval above 10ms and
below 30s, and every production caller takes the default.
TransactionAcquire::map_ and InboundLedger's trigger(), done() and
TriggerReason are protected, which is the only way a subclass reaches
them.

ConsensusTransSetSF::kMinTxNodeBytesToParse names the resubmission floor
as the 4-byte hash prefix plus kMinShaMapItemBytes plus one, which is
the same 17 the old size() > 16 test applied.
2026-10-04 06:35:35 -04:00
Bart
db21f68086 refactor: Expose a SHAMapAddNode verdict as counts
getBad() and getDuplicate() join getGood(), so a caller reads a verdict
as three counts rather than parsing the log string get() builds. The bad
count is per node rejected, not per node offered: a batch that stops at
its first rejected node counts one, and a batch that carries on counts
each. get() is a log format, so its wording is pinned in one place, the
new gtest, and nothing else depends on it.
2026-10-04 06:35:34 -04:00
10 changed files with 343 additions and 67 deletions

View File

@@ -14,35 +14,112 @@ private:
public:
SHAMapAddNode();
/**
* Record one node that was rejected.
*/
void
incInvalid();
/**
* Record one node that produced a good result.
*/
void
incUseful();
/**
* Record one node that was not needed: the map already held it, or the map
* or its acquisition was no longer taking nodes. isGood() counts it on the
* accepted side.
*/
void
incDuplicate();
void
reset();
/**
* @return How many nodes produced a good result.
*/
[[nodiscard]] int
getGood() const;
/**
* @return How many nodes this code rejected, which is not how many the
* batch was given: a batch that stops at its first bad node counts
* one, and a batch that carries on counts each.
*/
[[nodiscard]] int
getBad() const;
/**
* @return How many nodes were not needed, whether already held or offered
* to a map no longer taking nodes, tallied apart from the good and
* bad counts.
*/
[[nodiscard]] int
getDuplicate() const;
/**
* Whether nodes that produced a good result or were not needed outnumber
* the ones rejected.
*
* @return Whether the batch was good.
*/
[[nodiscard]] bool
isGood() const;
/**
* @return Whether at least one node in the batch was rejected.
*/
[[nodiscard]] bool
isInvalid() const;
/**
* @return Whether at least one node in the batch produced a good result.
*/
[[nodiscard]] bool
isUseful() const;
/**
* @return A verdict recording one node that was not needed.
*/
static SHAMapAddNode
duplicate();
/**
* @return A verdict recording one useful node.
*/
static SHAMapAddNode
useful();
/**
* @return A verdict recording one invalid node.
*/
static SHAMapAddNode
invalid();
/**
* Clear every count back to zero.
*/
void
reset();
/**
* Render the tally as a log line.
*
* @return The tally, e.g. "good:2 bad:1 dupe:1", or "no nodes processed" if
* every count is zero.
*/
[[nodiscard]] std::string
get() const;
/**
* Add another verdict's counts into this one.
*
* @param n The verdict to add.
* @return This verdict, updated.
*/
SHAMapAddNode&
operator+=(SHAMapAddNode const& n);
static SHAMapAddNode
duplicate();
static SHAMapAddNode
useful();
static SHAMapAddNode
invalid();
private:
SHAMapAddNode(int good, int bad, int duplicate);
};
@@ -74,18 +151,30 @@ SHAMapAddNode::incDuplicate()
++duplicate_;
}
inline void
SHAMapAddNode::reset()
{
good_ = bad_ = duplicate_ = 0;
}
inline int
SHAMapAddNode::getGood() const
{
return good_;
}
inline int
SHAMapAddNode::getBad() const
{
return bad_;
}
inline int
SHAMapAddNode::getDuplicate() const
{
return duplicate_;
}
inline bool
SHAMapAddNode::isGood() const
{
return (good_ + duplicate_) > bad_;
}
inline bool
SHAMapAddNode::isInvalid() const
{
@@ -98,22 +187,6 @@ SHAMapAddNode::isUseful() const
return good_ > 0;
}
inline SHAMapAddNode&
SHAMapAddNode::operator+=(SHAMapAddNode const& n)
{
good_ += n.good_;
bad_ += n.bad_;
duplicate_ += n.duplicate_;
return *this;
}
inline bool
SHAMapAddNode::isGood() const
{
return (good_ + duplicate_) > bad_;
}
inline SHAMapAddNode
SHAMapAddNode::duplicate()
{
@@ -132,6 +205,12 @@ SHAMapAddNode::invalid()
return SHAMapAddNode(0, 1, 0);
}
inline void
SHAMapAddNode::reset()
{
good_ = bad_ = duplicate_ = 0;
}
inline std::string
SHAMapAddNode::get() const
{
@@ -160,4 +239,14 @@ SHAMapAddNode::get() const
return ret;
}
inline SHAMapAddNode&
SHAMapAddNode::operator+=(SHAMapAddNode const& n)
{
good_ += n.good_;
bad_ += n.bad_;
duplicate_ += n.duplicate_;
return *this;
}
} // namespace xrpl

View File

@@ -0,0 +1,76 @@
#include <xrpl/shamap/SHAMapAddNode.h>
#include <gtest/gtest.h>
namespace xrpl::tests {
// get() is a log format, so its wording is pinned here, once. Every other site reads the same tally
// by value, through the count and verdict accessors.
TEST(SHAMapAddNode, get_names_every_non_empty_count)
{
EXPECT_EQ(SHAMapAddNode{}.get(), "no nodes processed");
EXPECT_EQ(SHAMapAddNode::useful().get(), "good:1");
EXPECT_EQ(SHAMapAddNode::invalid().get(), "bad:1");
EXPECT_EQ(SHAMapAddNode::duplicate().get(), "dupe:1");
// Several of a kind are counted, and the counts are joined in a fixed order with a single
// space, whichever order they were recorded in.
SHAMapAddNode san;
san.incInvalid();
san.incUseful();
san.incUseful();
san.incDuplicate();
EXPECT_EQ(san.get(), "good:2 bad:1 dupe:1");
san.reset();
EXPECT_EQ(san.get(), "no nodes processed");
}
// The three counts and the verdicts derived from them.
TEST(SHAMapAddNode, counts_and_verdicts_agree)
{
SHAMapAddNode san;
EXPECT_EQ(san.getGood(), 0);
EXPECT_EQ(san.getBad(), 0);
EXPECT_EQ(san.getDuplicate(), 0);
EXPECT_FALSE(san.isInvalid());
EXPECT_FALSE(san.isUseful());
// Good counts what produced a good result, and useful is that count being non-zero.
san.incUseful();
EXPECT_EQ(san.getGood(), 1);
EXPECT_TRUE(san.isUseful());
EXPECT_TRUE(san.isGood());
// A duplicate counts toward good. isUseful() here reflects the incUseful() above.
san.incDuplicate();
EXPECT_EQ(san.getDuplicate(), 1);
EXPECT_FALSE(san.isInvalid());
EXPECT_TRUE(san.isGood());
// Bad is a count, so a batch that carries on past a rejected node reports one per node, which
// distinguishes "stopped on the first" from "rejected several".
san.incInvalid();
EXPECT_EQ(san.getBad(), 1);
EXPECT_TRUE(san.isInvalid());
EXPECT_TRUE(san.isGood()) << "one bad node among two accepted ones is still a good batch";
san.incInvalid();
san.incInvalid();
EXPECT_EQ(san.getGood(), 1);
EXPECT_EQ(san.getBad(), 3);
EXPECT_EQ(san.getDuplicate(), 1);
EXPECT_FALSE(san.isGood()) << "more bad nodes than accepted ones is not a good batch";
// Adding one verdict to another sums every count.
SHAMapAddNode total;
total += SHAMapAddNode::useful();
total += SHAMapAddNode::invalid();
total += SHAMapAddNode::invalid();
total += SHAMapAddNode::duplicate();
EXPECT_EQ(total.getGood(), 1);
EXPECT_EQ(total.getBad(), 2);
EXPECT_EQ(total.getDuplicate(), 1);
}
} // namespace xrpl::tests

View File

@@ -180,4 +180,50 @@ TEST_F(SHAMapSyncTest, sync)
destination.invariants();
}
// The duplicate verdict also answers for a node that did not need to be added, which is the half
// of incDuplicate()'s contract that "already held" does not cover. addKnownNode() reaches it when
// the map has stopped taking nodes, and returns there before looking at the offer at all.
//
// An A/B on one offer and two maps of the same shape, so the synching state is the only difference
// between the two verdicts. Neither map holds the offered node, so a count that meant "already
// held" would be wrong for it.
TEST_F(SHAMapSyncTest, add_known_node_reports_duplicate_once_a_map_stops_synching)
{
TestNodeFamily f{j_};
// A well-formed inner node with one child. Its contents do not matter: both verdicts below
// are reached without the map reading them.
auto const makeOffer = [] {
Serializer s;
s.addBitString(UInt256{1});
s.add8(0);
s.add8(kWireTypeCompressedInner);
return SHAMapTreeNode::makeFromWire(makeSlice(s.peekData()));
};
ASSERT_TRUE(makeOffer());
SHAMapNodeID const target{1, UInt256{}};
// Synching, so the offer is examined, and refused because the empty root has no branch to
// hook it onto. The point is that the map looked.
SHAMap synching{SHAMapType::FREE, f};
synching.setSynching();
auto const examined = synching.addKnownNode(target, makeOffer(), nullptr);
EXPECT_TRUE(examined.isInvalid());
EXPECT_EQ(examined.getDuplicate(), 0);
// The same offer to a map that has stopped taking nodes.
SHAMap stopped{SHAMapType::FREE, f};
stopped.setSynching();
stopped.clearSynching();
ASSERT_FALSE(stopped.isSynching());
auto const notNeeded = stopped.addKnownNode(target, makeOffer(), nullptr);
EXPECT_EQ(notNeeded.getGood(), 0);
EXPECT_EQ(notNeeded.getBad(), 0);
EXPECT_EQ(notNeeded.getDuplicate(), 1);
EXPECT_TRUE(notNeeded.isGood()) << "a node that was not needed counts on the accepted side";
}
} // namespace xrpl::tests

View File

@@ -42,7 +42,7 @@ ConsensusTransSetSF::gotNode(
nodeCache_.insert(nodeHash, nodeData);
if ((type == SHAMapNodeType::TnTransactionNm) && (nodeData.size() > 16))
if ((type == SHAMapNodeType::TnTransactionNm) && (nodeData.size() >= kMinTxNodeBytesToParse))
{
// this is a transaction, and we didn't have it
JLOG(j_.debug()) << "Node on our acquiring TX set is TXN we may not have";

View File

@@ -9,6 +9,7 @@
#include <xrpl/shamap/SHAMapSyncFilter.h>
#include <xrpl/shamap/SHAMapTreeNode.h>
#include <cstddef>
#include <cstdint>
#include <optional>
@@ -24,6 +25,15 @@ class ConsensusTransSetSF : public SHAMapSyncFilter
public:
using NodeCache = TaggedCache<SHAMapHash, Blob>;
/**
* The size a node's hash-prefixed wire data must reach before gotNode()
* tries to parse and resubmit it as a transaction. One byte past the
* smallest a hash-prefixed SHAMap leaf can be, which a signed transaction
* clears.
*/
static constexpr std::size_t kMinTxNodeBytesToParse =
sizeof(std::uint32_t) + kMinShaMapItemBytes + 1;
ConsensusTransSetSF(Application& app, NodeCache& nodeCache);
// Note that the nodeData is overwritten by this call

View File

@@ -30,9 +30,9 @@
namespace xrpl {
// A ledger we are trying to acquire
class InboundLedger final : public TimeoutCounter,
public std::enable_shared_from_this<InboundLedger>,
public CountedObject<InboundLedger>
class InboundLedger : public TimeoutCounter,
public std::enable_shared_from_this<InboundLedger>,
public CountedObject<InboundLedger>
{
public:
using ClockType = beast::AbstractClock<std::chrono::steady_clock>;
@@ -44,13 +44,29 @@ public:
CONSENSUS // We believe the consensus round requires this ledger
};
/**
* How long to wait between retries, and so how long one timeout takes.
*/
static constexpr std::chrono::milliseconds kRetryInterval{3000};
/**
* @param app The application to run in.
* @param hash The ledger to acquire.
* @param seq Its sequence, or zero if not known yet.
* @param reason Why it is being acquired.
* @param clock The clock touch() records against.
* @param peerSet Which peers to ask, and how to reach them.
* @param retryInterval How long to wait between retries. TimeoutCounter
* requires more than 10ms and less than 30s.
*/
InboundLedger(
Application& app,
UInt256 const& hash,
std::uint32_t seq,
Reason reason,
ClockType&,
std::unique_ptr<PeerSet> peerSet);
ClockType& clock,
std::unique_ptr<PeerSet> peerSet,
std::chrono::milliseconds retryInterval = kRetryInterval);
~InboundLedger() override;
@@ -119,15 +135,33 @@ public:
return lastAction_;
}
private:
protected:
// Protected so a test subclass can drive an acquisition through trigger() and done().
// Why trigger() is being run, which decides how deep a request goes and whether an
// aggressive retry applies.
enum class TriggerReason { Added, Reply, Timeout };
/**
* Ask for more nodes, or judge what has been collected.
*
* @param peer The peer to ask, or nullptr to ask everyone being tracked.
* @param reason Why the acquisition is being triggered.
*/
void
trigger(std::shared_ptr<Peer> const& peer, TriggerReason reason);
/**
* Settle the acquisition and signal whatever is waiting on it. Runs at most
* once. Callers hold mtx_, except trigger(), which releases it first.
*/
void
done();
private:
void
filterNodes(std::vector<std::pair<SHAMapNodeID, UInt256>>& nodes, TriggerReason reason);
void
trigger(std::shared_ptr<Peer> const&, TriggerReason);
std::vector<NeededHashT>
getNeededHashes();
@@ -138,10 +172,7 @@ private:
tryDB(node_store::Database& srcDB);
void
done();
void
onTimer(bool progress, ScopedLockType& peerSetLock) override;
onTimer(bool progress, ScopedLockType& sl) override;
std::size_t
getPeerCount() const;

View File

@@ -54,8 +54,6 @@
namespace xrpl {
using namespace std::chrono_literals;
static constexpr auto kPeerCountStart = 5; // Number of peers to start with
static constexpr auto kPeerCountAdd = 3; // Number of peers to add on a timeout
static constexpr auto kLedgerTimeoutRetriesMax = 6; // how many timeouts before we give up
@@ -65,20 +63,18 @@ static constexpr auto kMissingNodesFind = 256; // Number of nodes to find initi
static constexpr auto kReqNodesReply = 128; // Number of nodes to request for a reply
static constexpr auto kReqNodes = 12; // Number of nodes to request blindly
// millisecond for each ledger timeout
constexpr auto kLedgerAcquireTimeout = 3000ms;
InboundLedger::InboundLedger(
Application& app,
UInt256 const& hash,
std::uint32_t seq,
Reason reason,
ClockType& clock,
std::unique_ptr<PeerSet> peerSet)
std::unique_ptr<PeerSet> peerSet,
std::chrono::milliseconds retryInterval)
: TimeoutCounter(
app,
hash,
kLedgerAcquireTimeout,
retryInterval,
{.jobType = JtLedgerData, .jobName = "InboundLedger", .jobLimit = 5},
app.getJournal("InboundLedger"))
, clock_(clock)

View File

@@ -98,9 +98,14 @@ protected:
/**
* Hook called from invokeOnTimer().
*
* @param progress Whether the subtype recorded progress since the
* last call.
* @param sl Proof mtx_ is held. mtx_ is recursive, so a nested lock taken
* inside the call does not release it.
*/
virtual void
onTimer(bool progress, ScopedLockType&) = 0;
onTimer(bool progress, ScopedLockType& sl) = 0;
/**
* Return a weak pointer to this.

View File

@@ -18,6 +18,7 @@
#include <xrpl.pb.h>
#include <algorithm>
#include <chrono>
#include <cstddef>
#include <exception>
#include <memory>
@@ -26,22 +27,18 @@
namespace xrpl {
using namespace std::chrono_literals;
// Timeout interval in milliseconds
constexpr auto kTxAcquireTimeout = 250ms;
static constexpr auto kNormTimeouts = 4;
static constexpr auto kMaxTimeouts = 20;
TransactionAcquire::TransactionAcquire(
Application& app,
UInt256 const& hash,
std::unique_ptr<PeerSet> peerSet)
std::unique_ptr<PeerSet> peerSet,
std::chrono::milliseconds retryInterval)
: TimeoutCounter(
app,
hash,
kTxAcquireTimeout,
retryInterval,
{.jobType = JtTxnData, .jobName = "TxAcq", .jobLimit = {}},
app.getJournal("TransactionAcquire"))
, peerSet_(std::move(peerSet))
@@ -79,7 +76,7 @@ TransactionAcquire::done()
}
void
TransactionAcquire::onTimer(bool progress, ScopedLockType& psl)
TransactionAcquire::onTimer(bool progress, ScopedLockType&)
{
if (timeouts_ > kMaxTimeouts)
{

View File

@@ -12,6 +12,7 @@
#include <xrpl/shamap/SHAMapNodeID.h>
#include <xrpl/shamap/SHAMapTreeNode.h>
#include <chrono>
#include <cstddef>
#include <memory>
#include <utility>
@@ -21,14 +22,30 @@ namespace xrpl {
// VFALCO TODO rename to PeerTxRequest
// A transaction set we are trying to acquire
class TransactionAcquire final : public TimeoutCounter,
public std::enable_shared_from_this<TransactionAcquire>,
public CountedObject<TransactionAcquire>
class TransactionAcquire : public TimeoutCounter,
public std::enable_shared_from_this<TransactionAcquire>,
public CountedObject<TransactionAcquire>
{
public:
using pointer = std::shared_ptr<TransactionAcquire>;
TransactionAcquire(Application& app, UInt256 const& hash, std::unique_ptr<PeerSet> peerSet);
/**
* How long to wait between retries, and so how long one timeout takes.
*/
static constexpr std::chrono::milliseconds kRetryInterval{250};
/**
* @param app The application to run in.
* @param hash The set to acquire.
* @param peerSet Which peers to ask, and how to reach them.
* @param retryInterval How long to wait between retries. TimeoutCounter
* requires more than 10ms and less than 30s.
*/
TransactionAcquire(
Application& app,
UInt256 const& hash,
std::unique_ptr<PeerSet> peerSet,
std::chrono::milliseconds retryInterval = kRetryInterval);
~TransactionAcquire() override = default;
SHAMapAddNode
@@ -42,14 +59,23 @@ public:
void
stillNeed();
private:
protected:
// Kept protected so a test subclass can read the map's state.
// Production callers reach a set through InboundTransactions.
std::shared_ptr<SHAMap> map_;
private:
bool haveRoot_{false};
std::unique_ptr<PeerSet> peerSet_;
void
onTimer(bool progress, ScopedLockType& peerSetLock) override;
onTimer(bool progress, ScopedLockType& sl) override;
/**
* Settle the acquired set and hand it on, or report the failure. Call under
* mtx_. Runs once per outcome, and a stillNeed() revival gives one object a
* second outcome.
*/
void
done();