Compare commits

...

7 Commits

Author SHA1 Message Date
Bart
ccd8838d09 fix: Report a node the sync path refuses as invalid data
addKnownNode() returns invalid() for the two shapes it declines to hook
in: an inner node at kLeafDepth, a depth only a leaf may occupy, and a
node whose claimed ID is not where the hash-verified descent stopped.
The depth case also marks the map Invalid, which is terminal, since
every node above it hash-verified and so the requested root hash itself
commits to a shape no valid tree has. An ID mismatch leaves the map
sound. The depth test precedes the full-below cache lookup, which is
keyed by hash and so answers for no particular depth.

addRootNode() guarantees in every build mode that a root already held is
a duplicate only if it hashes to the hash the map was asked for, and
invalid data otherwise.
2026-10-04 06:35:38 -04:00
Bart
dcd44fa0ae fix: Make the SHAMap sync-path state atomic
SHAMap::state_, ::full_, ::ledgerSeq_ and SHAMapInnerNode::fullBelowGen_
are std::atomic and pinned lock-free, since a nodestore fetch and a
getMissingNodes() walk reach them while the thread driving an
acquisition does, and neither may block. finishFetch() withdraws full_
with an acq_rel exchange, so exactly one reader reports the gap.
ledgerSeq_ and fullBelowGen_ are relaxed both ways: a lookup hint for a
store keyed by hash, and a generation compared only for equality whose
children are published through the node's own lock. Ledger::setFull()
stores each map's sequence before its full flag, the one order that
makes the sequence visible to the thread whose exchange wins.

The shamap tests' TestNodeFamily replaces helpers/TestFamily.h, takes a
readThreads count for its nodestore, and declares its clock ahead of the
members that use it. InnerNode.h adds makeFullInnerNode() and
makeCompressedInnerNode(), which build an inner node from a list of
children.
2026-10-04 06:35:37 -04:00
Bart
a0fd44cc7b fix: Judge a header's account hash on both acquisition routes
No ledger has an empty state map, so a header with a zero account hash
cannot name a ledger. tryDB() judges that inside makeLedger(), next to
the hash and sequence check that already refuses a header that cannot be
a ledger, and takeHeader() judges it through the same helper,
failOnZeroAccountHash(). The helper sets failed_ and drops ledger_, so
neither route stores the header, sets seq_, or sets haveHeader_ for such
a header. A zero transaction hash is still an empty transaction set.

The header is the one asked for, so the peer is not charged.
processData() calls done() when takeHeader() has failed the acquisition,
since trigger() and onTimer() return early once isDone() and nothing else
would signal the waiters or record the hash in recentFailures_.
2026-10-04 06:35:37 -04:00
Bart
516ed4d433 fix: Signal every InboundLedger failure found in local data
init() and trigger() call done() when tryDB() sets failed_, so every
route that ends an acquisition runs the one function that signals
whatever waits on it and records the hash in recentFailures_.
checkLocal() already did. A hash in recentFailures_ is what keeps a
later round from asking for a ledger already judged unobtainable.
2026-10-04 06:35:36 -04:00
Bart
a790fb3527 refactor: Add a reusable peer harness for acquisition tests
DeepChain builds the node chains both acquisition suites need: chains
that run to SHAMap::kLeafDepth, which no valid tree holds, and toLeaf()
chains that complete an acquisition. AcquireTestHelpers.h adds
ChargeRecordingPeer, RequestCountingPeerSet, packetFor(), waitFor() and
tallyIs(), so each suite drives an acquisition through the real
gotData() dispatch rather than reproducing it.

AcquireTestHelpers.h is the first src/test file to include one from
src/tests, so ordering.txt records a test.app > tests.libxrpl edge. That
direction is deliberate: src/test goes away as its suites migrate, so an
edge into the surviving tree outlives the migration. Nothing under
src/tests includes src/test.
2026-10-04 06:35:36 -04:00
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
26 changed files with 2815 additions and 269 deletions

View File

@@ -103,6 +103,7 @@ words:
- dearmor
- decryptor
- dedented
- dedup
- deleteme
- demultiplexer
- deserializaton
@@ -344,6 +345,7 @@ words:
- unambiguity
- unauthorizes
- unauthorizing
- undeserializable
- unergonomic
- unfetched
- unfindable

View File

@@ -62,6 +62,7 @@ libxrpl.tx > xrpl.protocol
libxrpl.tx > xrpl.server
libxrpl.tx > xrpl.tx
test.app > test.jtx
test.app > tests.libxrpl
test.app > test.unit_test
test.app > xrpl.basics
test.app > xrpl.config

View File

@@ -285,10 +285,13 @@ public:
void
setFull() const
{
txMap_.setFull();
// Sequence before flag, per map: setLedgerSeq() stores relaxed and setFull() stores
// release, and SHAMap::finishFetch() reads the sequence only after acquiring the flag, so
// only this order publishes it.
txMap_.setLedgerSeq(header_.seq);
stateMap_.setFull();
txMap_.setFull();
stateMap_.setLedgerSeq(header_.seq);
stateMap_.setFull();
}
void

View File

@@ -17,6 +17,7 @@
#include <xrpl/shamap/SHAMapNodeID.h>
#include <xrpl/shamap/SHAMapTreeNode.h>
#include <atomic>
#include <condition_variable>
#include <cstddef>
#include <cstdint>
@@ -40,7 +41,7 @@ class SHAMapSyncFilter;
/**
* Describes the current state of a given SHAMap
*/
enum class SHAMapState {
enum class SHAMapState : std::uint8_t {
/**
* The map is in flux and objects can be added and removed.
*
@@ -120,16 +121,41 @@ private:
*/
std::uint32_t cowid_ = 1;
// ledgerSeq_, state_ and full_ are touched on the nodestore fetch path and on a
// getMissingNodes() walk, neither of which may block, so pin them lock-free.
static_assert(std::atomic<std::uint32_t>::is_always_lock_free);
static_assert(std::atomic<SHAMapState>::is_always_lock_free);
static_assert(std::atomic<bool>::is_always_lock_free);
/**
* The sequence of the ledger that this map references, if any.
*
* Written when a map's ledger sequence is established (Ledger::setFull(),
* InboundLedger) while a nodestore reader thread reads it. Relaxed either
* way: it serves as a lookup hint for a nodestore keyed by hash.
*/
std::uint32_t ledgerSeq_ = 0;
std::atomic<std::uint32_t> ledgerSeq_ = 0;
SHAMapTreeNodePtr root_;
mutable SHAMapState state_;
/**
* The map's state.
*
* A getMissingNodes() walk writes it, through clearSynching(), while whatever
* drives the acquisition reads it. Atomic rather than guarded, since the
* acquisition code releases its lock across that walk.
*/
mutable std::atomic<SHAMapState> state_;
SHAMapType const type_;
bool backed_ = true; // Map is backed by the database
mutable bool full_ = false; // Map is believed complete in database
bool backed_ = true; // Map is backed by the database
/**
* Map is believed complete in database.
*
* Atomic because finishFetch() clears it on any nodestore reader thread,
* several of which can run at once.
*/
mutable std::atomic<bool> full_ = false;
public:
/**
@@ -341,6 +367,9 @@ public:
* This function is used when receiving the root node of a SHAMap from a peer during ledger
* synchronization. The node must already have been deserialized.
*
* A root offered under a hash the map does not hold names another tree and
* is reported as invalid data. A root matching that hash is a duplicate.
*
* @param hash The expected hash of the root node.
* @param rootNode A deserialized root node to add.
* @param filter Optional sync filter to track received nodes.
@@ -360,6 +389,9 @@ public:
* is inserted at the position specified by nodeID. The node must already have been
* deserialized.
*
* A node that no valid tree can hold makes the map Invalid: the root hash
* committed to an impossible shape, so no peer can satisfy it.
*
* @param nodeID The position in the tree where this node belongs.
* @param treeNode A deserialized tree node to add.
* @param filter Optional sync filter to track received nodes.
@@ -375,16 +407,29 @@ public:
SHAMapTreeNodePtr treeNode,
SHAMapSyncFilter const* filter);
// status functions
void
setImmutable();
bool
/**
* @return Whether the map is being synced against a hash it was given, so
* its hash is fixed while nodes may still be added.
*/
[[nodiscard]] bool
isSynching() const;
void
setSynching();
void
clearSynching();
bool
/**
* Whether the map can still be the map it claims to be.
*
* A map that is merely missing nodes is valid, and stays valid until
* something proves the tree it is syncing against cannot exist.
*
* @return Whether the map has not been proven impossible.
*/
[[nodiscard]] bool
isValid() const;
// caution: otherMap must be accessed only by this function
@@ -524,6 +569,34 @@ private:
using DeltaRef =
std::pair<boost::intrusive_ptr<SHAMapItem const>, boost::intrusive_ptr<SHAMapItem const>>;
/**
* The sequence of the ledger this map references, read atomically.
*
* @return The sequence, or zero if the map references no ledger.
*/
[[nodiscard]] std::uint32_t
ledgerSeq() const;
/**
* The current state, read atomically.
*
* Orders state_ alone. The tree's nodes carry no ordering guarantees.
*
* @return The state as of the call, which a concurrent walk may already
* have moved past.
*/
[[nodiscard]] SHAMapState
state() const;
/**
* Record that the map is provably not the one it claims to be.
*
* Private because only the map itself can prove that, from a node that
* contradicts the hashes it is syncing against.
*/
void
setInvalid();
// tree node cache operations
SHAMapTreeNodePtr
cacheLookup(SHAMapHash const& hash) const;
@@ -741,44 +814,62 @@ private:
inline void
SHAMap::setFull()
{
full_ = true;
full_.store(true, std::memory_order_release);
}
inline void
SHAMap::setLedgerSeq(std::uint32_t lseq)
{
ledgerSeq_ = lseq;
ledgerSeq_.store(lseq, std::memory_order_relaxed);
}
inline std::uint32_t
SHAMap::ledgerSeq() const
{
return ledgerSeq_.load(std::memory_order_relaxed);
}
inline SHAMapState
SHAMap::state() const
{
return state_.load(std::memory_order_acquire);
}
inline void
SHAMap::setImmutable()
{
XRPL_ASSERT(state_ != SHAMapState::Invalid, "xrpl::SHAMap::setImmutable : state is valid");
state_ = SHAMapState::Immutable;
XRPL_ASSERT(isValid(), "xrpl::SHAMap::setImmutable : state is valid");
state_.store(SHAMapState::Immutable, std::memory_order_release);
}
inline bool
SHAMap::isSynching() const
{
return state_ == SHAMapState::Synching;
return state() == SHAMapState::Synching;
}
inline void
SHAMap::setSynching()
{
state_ = SHAMapState::Synching;
state_.store(SHAMapState::Synching, std::memory_order_release);
}
inline void
SHAMap::clearSynching()
{
state_ = SHAMapState::Modifying;
state_.store(SHAMapState::Modifying, std::memory_order_release);
}
inline bool
SHAMap::isValid() const
{
return state_ != SHAMapState::Invalid;
return state() != SHAMapState::Invalid;
}
inline void
SHAMap::setInvalid()
{
state_.store(SHAMapState::Invalid, std::memory_order_release);
}
inline void

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

@@ -31,7 +31,19 @@ private:
*/
TaggedPointer hashesAndChildren_;
std::uint32_t fullBelowGen_ = 0;
// Pin that wrapping fullBelowGen_ in an atomic keeps the packed layout, and that the load
// isFullBelow() does once per node of every walk stays lock-free.
static_assert(std::atomic<std::uint32_t>::is_always_lock_free);
static_assert(sizeof(std::atomic<std::uint32_t>) == sizeof(std::uint32_t));
static_assert(alignof(std::atomic<std::uint32_t>) == alignof(std::uint32_t));
/**
* Written from more than one thread, since canonicalization shares a node
* between maps and a walk can run with the acquisition lock released (see
* SHAMap::state_). Relaxed both ways, since a generation is only compared
* for equality and the children it vouches for are published through lock_.
*/
std::atomic<std::uint32_t> fullBelowGen_ = 0;
std::uint16_t isBranch_ = 0;
/**
@@ -204,13 +216,13 @@ SHAMapInnerNode::getBranchCount() const
inline bool
SHAMapInnerNode::isFullBelow(std::uint32_t generation) const
{
return fullBelowGen_ == generation;
return fullBelowGen_.load(std::memory_order_relaxed) == generation;
}
inline void
SHAMapInnerNode::setFullBelowGen(std::uint32_t gen)
{
fullBelowGen_ = gen;
fullBelowGen_.store(gen, std::memory_order_relaxed);
}
} // namespace xrpl

View File

@@ -26,6 +26,7 @@
#include <boost/smart_ptr/intrusive_ptr.hpp>
#include <atomic>
#include <cstdint>
#include <exception>
#include <functional>
@@ -77,14 +78,14 @@ SHAMap::SHAMap(SHAMap const& other, bool isMutable)
: f_(other.f_)
, journal_(other.f_.journal())
, cowid_(other.cowid_ + 1)
, ledgerSeq_(other.ledgerSeq_)
, ledgerSeq_(other.ledgerSeq())
, root_(other.root_)
, state_(isMutable ? SHAMapState::Modifying : SHAMapState::Immutable)
, type_(other.type_)
, backed_(other.backed_)
{
// If either map may change, they cannot share nodes
if ((state_ != SHAMapState::Immutable) || (other.state_ != SHAMapState::Immutable))
if ((state() != SHAMapState::Immutable) || (other.state() != SHAMapState::Immutable))
{
unshare();
}
@@ -105,7 +106,7 @@ SHAMap::dirtyUp(NodePathStack& stack, UInt256 const& target, SHAMapTreeNodePtr c
// child can be an inner node or a leaf
XRPL_ASSERT(
(state_ != SHAMapState::Synching) && (state_ != SHAMapState::Immutable),
(state() != SHAMapState::Synching) && (state() != SHAMapState::Immutable),
"xrpl::SHAMap::dirtyUp : valid state");
XRPL_ASSERT(child && (child->cowid() == cowid_), "xrpl::SHAMap::dirtyUp : valid child input");
@@ -170,7 +171,7 @@ SHAMapTreeNodePtr
SHAMap::fetchNodeFromDB(SHAMapHash const& hash) const
{
XRPL_ASSERT(backed_, "xrpl::SHAMap::fetchNodeFromDB : is backed");
auto obj = f_.db().fetchNodeObject(hash.asUInt256(), ledgerSeq_);
auto obj = f_.db().fetchNodeObject(hash.asUInt256(), ledgerSeq());
return finishFetch(hash, obj);
}
@@ -183,10 +184,14 @@ SHAMap::finishFetch(SHAMapHash const& hash, std::shared_ptr<NodeObject> const& o
{
if (!object)
{
if (full_)
// A missing node disproves full_, so withdraw it and report the gap. The acq_rel
// exchange orders the report and makes exactly one reader report it. The load ahead
// of it is relaxed, since all it decides is whether to try the exchange: a stale true
// costs one exchange that loses, and a stale false means another reader won already.
if (full_.load(std::memory_order_relaxed) &&
full_.exchange(false, std::memory_order_acq_rel))
{
full_ = false;
f_.missingNodeAcquireBySeq(ledgerSeq_, hash.asUInt256());
f_.missingNodeAcquireBySeq(ledgerSeq(), hash.asUInt256());
}
return {};
}
@@ -219,7 +224,7 @@ SHAMap::checkFilter(SHAMapHash const& hash, SHAMapSyncFilter const* filter) cons
auto node = SHAMapTreeNode::makeFromPrefix(makeSlice(*nodeData), hash);
if (node)
{
filter->gotNode(true, hash, ledgerSeq_, std::move(*nodeData), node->getType());
filter->gotNode(true, hash, ledgerSeq(), std::move(*nodeData), node->getType());
if (backed_)
canonicalize(hash, node);
}
@@ -399,7 +404,7 @@ SHAMap::descendAsync(
{
f_.db().asyncFetch(
hash.asUInt256(),
ledgerSeq_,
ledgerSeq(),
[this, hash, cb{std::move(callback)}](std::shared_ptr<NodeObject> const& object) {
auto node = finishFetch(hash, object);
cb(node, hash);
@@ -424,7 +429,7 @@ SHAMap::unshareNode(intr_ptr::SharedPtr<Node> node, SHAMapNodeID const& nodeID)
if (node->cowid() != cowid_)
{
// have a CoW
XRPL_ASSERT(state_ != SHAMapState::Immutable, "xrpl::SHAMap::unshareNode : not immutable");
XRPL_ASSERT(state() != SHAMapState::Immutable, "xrpl::SHAMap::unshareNode : not immutable");
node = intr_ptr::staticPointerCast<Node>(node->clone(cowid_));
if (nodeID.isRoot())
root_ = node;
@@ -637,7 +642,7 @@ bool
SHAMap::delItem(UInt256 const& id)
{
// delete the item with this ID
XRPL_ASSERT(state_ != SHAMapState::Immutable, "xrpl::SHAMap::delItem : not immutable");
XRPL_ASSERT(state() != SHAMapState::Immutable, "xrpl::SHAMap::delItem : not immutable");
NodePathStack stack;
walkTowardsKey(id, &stack);
@@ -719,7 +724,7 @@ SHAMap::delItem(UInt256 const& id)
bool
SHAMap::addGiveItem(SHAMapNodeType type, boost::intrusive_ptr<SHAMapItem const> item)
{
XRPL_ASSERT(state_ != SHAMapState::Immutable, "xrpl::SHAMap::addGiveItem : not immutable");
XRPL_ASSERT(state() != SHAMapState::Immutable, "xrpl::SHAMap::addGiveItem : not immutable");
XRPL_ASSERT(type != SHAMapNodeType::TnInner, "xrpl::SHAMap::addGiveItem : valid type input");
// add the specified item, does not update
@@ -810,7 +815,7 @@ SHAMap::updateGiveItem(SHAMapNodeType type, boost::intrusive_ptr<SHAMapItem cons
// can't change the tag but can change the hash
UInt256 const tag = item->key();
XRPL_ASSERT(state_ != SHAMapState::Immutable, "xrpl::SHAMap::updateGiveItem : not immutable");
XRPL_ASSERT(state() != SHAMapState::Immutable, "xrpl::SHAMap::updateGiveItem : not immutable");
NodePathStack stack;
walkTowardsKey(tag, &stack);
@@ -901,7 +906,7 @@ SHAMap::writeNode(NodeObjectType t, SHAMapTreeNodePtr node) const
Serializer s;
node->serializeWithPrefix(s);
f_.db().store(t, std::move(s.modData()), node->getHash().asUInt256(), ledgerSeq_);
f_.db().store(t, std::move(s.modData()), node->getHash().asUInt256(), ledgerSeq());
return node;
}

View File

@@ -16,6 +16,7 @@
#include <xrpl/shamap/detail/TaggedPointer.h>
#include <xrpl/shamap/detail/TaggedPointer.ipp>
#include <atomic>
#include <cstddef>
#include <cstdint>
#include <mutex>
@@ -77,7 +78,9 @@ SHAMapInnerNode::clone(std::uint32_t cowid) const
auto p = intr_ptr::makeShared<SHAMapInnerNode>(cowid, branchCount);
p->hash_ = hash_;
p->isBranch_ = isBranch_;
p->fullBelowGen_ = fullBelowGen_;
// Relaxed, as everywhere. p is not reachable by another thread until this returns.
p->fullBelowGen_.store(
fullBelowGen_.load(std::memory_order_relaxed), std::memory_order_relaxed);
SHAMapHash* cloneHashes = nullptr;
SHAMapHash* thisHashes = nullptr;
SHAMapTreeNodePtr* cloneChildren = nullptr;

View File

@@ -31,6 +31,25 @@
namespace xrpl {
namespace {
/**
* Whether a depth is at or past the deepest an inner node may occupy.
*
* Nibbles run out at SHAMap::kLeafDepth, so only a leaf may sit there. True for
* every deeper position too, which lies past the end of a key.
*
* @param depth The depth to judge.
* @return Whether an inner node at that depth makes the map impossible.
*/
[[nodiscard]] bool
isLeafDepth(unsigned int depth)
{
return depth >= SHAMap::kLeafDepth;
}
} // namespace
void
SHAMap::visitLeaves(
std::function<void(boost::intrusive_ptr<SHAMapItem const> const& item)> const& leafFunction)
@@ -527,11 +546,19 @@ SHAMap::addRootNode(
XRPL_ASSERT(cowid_ >= 1, "xrpl::SHAMap::addRootNode : valid cowid");
XRPL_ASSERT(rootNode, "xrpl::SHAMap::addRootNode : non-null root node");
// we already have a root_ node
// A map syncs against one hash and installs a root once, so a root already held is a duplicate
// only once it hashes to the hash asked for.
if (root_->getHash().isNonZero())
{
JLOG(journal_.trace()) << "Got root node, already have one";
XRPL_ASSERT(root_->getHash() == hash, "xrpl::SHAMap::addRootNode : valid hash");
if (root_->getHash() != hash)
{
JLOG(journal_.warn()) << "Root node offered under hash " << hash
<< ", but the map holds " << root_->getHash();
return SHAMapAddNode::invalid();
}
return SHAMapAddNode::duplicate();
}
@@ -555,7 +582,7 @@ SHAMap::addRootNode(
Serializer s;
root_->serializeWithPrefix(s);
filter->gotNode(
false, root_->getHash(), ledgerSeq_, std::move(s.modData()), root_->getType());
false, root_->getHash(), ledgerSeq(), std::move(s.modData()), root_->getType());
}
return SHAMapAddNode::useful();
@@ -598,7 +625,13 @@ SHAMap::addKnownNode(
}
auto childHash = inner->getChildHash(branch);
if (f_.getFullBelowCache()->touchIfExists(childHash.asUInt256()))
// Depth before the cache: the cache is keyed by node hash and shared across the family, and
// a hash covers a node's children but not its depth, so the same subtree can be cached as
// complete at one depth and reached at another. Skipping the shortcut only forgoes an
// optimization.
if (!isLeafDepth(currNodeID.getDepth() + 1) &&
f_.getFullBelowCache()->touchIfExists(childHash.asUInt256()))
{
return SHAMapAddNode::duplicate();
}
@@ -616,25 +649,32 @@ SHAMap::addKnownNode(
return SHAMapAddNode::invalid();
}
// Inner nodes must be at a level strictly less than 64
// but leaf nodes (while notionally at level 64) can be
// at any depth up to and including 64:
if ((currNodeID.getDepth() > kLeafDepth) ||
(treeNode->isInner() && currNodeID.getDepth() == kLeafDepth))
// Only leaves may sit at kLeafDepth (see isLeafDepth), so an inner node there makes the map
// impossible. Nothing is hooked in, so this is bad data rather than progress.
//
// Every node from the root down hash-verified to get here, so it is the requested root hash
// itself that commits to a shape no valid tree can have. The verdict belongs to that hash
// rather than to our copy of the tree: no peer can satisfy it, so retrying is futile.
bool const badDepth = treeNode->isInner() && isLeafDepth(currNodeID.getDepth());
SOMETIMES(badDepth, "xrpl::SHAMap::addKnownNode : map is invalid");
if (badDepth)
{
// Map is provably invalid
state_ = SHAMapState::Invalid;
return SHAMapAddNode::useful();
JLOG(journal_.warn()) << "Node " << nodeID << " makes the map invalid at "
<< currNodeID;
setInvalid();
return SHAMapAddNode::invalid();
}
if (currNodeID != nodeID)
// The data hashes to the child at currNodeID but claims to belong at nodeID, so it is not
// the node that was asked for. Only the label is wrong, so the map stays sound and the node
// is still obtainable from another sender.
bool const badPosition = (currNodeID != nodeID);
SOMETIMES(badPosition, "xrpl::SHAMap::addKnownNode : node ID does not match its position");
if (badPosition)
{
// Either this node is broken or we didn't request it (yet)
JLOG(journal_.warn()) << "unable to hook node " << nodeID;
JLOG(journal_.info()) << " stuck at " << currNodeID;
JLOG(journal_.info()) << "got depth=" << nodeID.getDepth()
<< ", walked to= " << currNodeID.getDepth();
return SHAMapAddNode::useful();
JLOG(journal_.warn()) << "Unable to hook node " << nodeID << ", stuck at "
<< currNodeID;
return SHAMapAddNode::invalid();
}
if (backed_)
@@ -647,7 +687,7 @@ SHAMap::addKnownNode(
Serializer s;
treeNode->serializeWithPrefix(s);
filter->gotNode(
false, childHash, ledgerSeq_, std::move(s.modData()), treeNode->getType());
false, childHash, ledgerSeq(), std::move(s.modData()), treeNode->getType());
}
return SHAMapAddNode::useful();
@@ -766,7 +806,7 @@ SHAMap::hasLeafNode(UInt256 const& tag, SHAMapHash const& targetNodeHash) const
// Same kLeafDepth hazard as in visitDifferences above. That guard bounds the caller's own
// traversal, not the map queried here, and the loop below descends from this map's root
// independently, so this check is what keeps a malformed map from reaching getChildNodeID.
if (nodeID.getDepth() >= kLeafDepth)
if (isLeafDepth(nodeID.getDepth()))
{
// LCOV_EXCL_START
UNREACHABLE("xrpl::SHAMap::hasLeafNode : inner node at leaf depth");

View File

@@ -0,0 +1,365 @@
#pragma once
#include <test/jtx/PeerStub.h>
#include <xrpld/app/ledger/ConsensusTransSetSF.h>
#include <xrpld/overlay/Peer.h>
#include <xrpld/overlay/PeerSet.h>
#include <xrpl/basics/SHAMapHash.h>
#include <xrpl/basics/base_uint.h>
#include <xrpl/protocol/Serializer.h>
#include <xrpl/resource/Charge.h>
#include <xrpl/shamap/SHAMapAddNode.h>
#include <xrpl/shamap/SHAMapNodeID.h>
#include <xrpl/shamap/SHAMapTreeNode.h>
#include <tests/libxrpl/shamap/DeepChain.h>
#include <xrpl.pb.h>
#include <atomic>
#include <chrono>
#include <cstddef>
#include <cstdint>
#include <functional>
#include <memory>
#include <mutex>
#include <optional>
#include <set>
#include <string>
#include <thread>
#include <utility>
#include <vector>
namespace xrpl::test {
// The chain builder needs only libxrpl, so it is shared with the gtest suites; see the header for
// what keeps the protobuf reply builder below in this tree.
using tests::DeepChain;
// A smallest-possible leaf plus its 4-byte HashPrefix lands one byte short of the floor
// ConsensusTransSetSF::gotNode() parses at, so a chain's leaf stays below the parse threshold.
static_assert(
sizeof(std::uint32_t) + DeepChain::kLeafItemBytes < ConsensusTransSetSF::kMinTxNodeBytesToParse,
"a smallest-possible leaf must stay below the resubmission floor");
/**
* A peer that records what it was charged. Every other method comes from
* PeerStub.
*
* One instance per packet keeps charges() unambiguous about which packet was
* charged what.
*
* charges_ is unguarded: charging happens on the packet path, so every charge
* lands on the thread that fed the packet in.
*/
class ChargeRecordingPeer : public PeerStub
{
public:
/**
* @param hasTxSet What hasTxSet() reports, which is how an acquisition
* decides whether this peer is worth asking. Defaults to true, so a
* peer handed straight to takeNodes() needs no argument.
*/
explicit ChargeRecordingPeer(bool hasTxSet = true) : PeerStub(nextId()), hasTxSet_(hasTxSet)
{
}
void
charge(resource::Charge const& fee, std::string const&) override
{
charges_.push_back(fee);
}
/**
* @return Every fee this peer has been charged, in the order charged.
*/
[[nodiscard]] std::vector<resource::Charge> const&
charges() const
{
return charges_;
}
// PeerStub returns false for both, and an acquisition asks the peers reporting true.
[[nodiscard]] bool
hasTxSet(UInt256 const&) const override
{
return hasTxSet_;
}
[[nodiscard]] bool
hasLedger(UInt256 const&, std::uint32_t) const override
{
return true;
}
private:
/**
* The next id to hand out, distinct per instance because
* RequestCountingPeerSet dedups by tracked id and PeerStub's own default is
* zero for every instance.
*
* @return The id.
*/
[[nodiscard]] static ID
nextId()
{
static std::atomic<ID> next{1};
return next++;
}
std::vector<resource::Charge> charges_;
bool hasTxSet_;
};
/**
* A peer set that counts the requests an acquisition makes through it. Offers
* peers to a hasItem/onPeerAdded callback pair, hard-filtered by hasItem (which
* only scores in the real peer set) and deduped by tracked id.
*
* A count is a call, not a delivery: a request naming no peer reaches nobody
* while this set tracks none. requests() and broadcasts() are counted apart.
*
* Every write here is guarded, since the retry timer drives addPeers() and
* sendRequest() from a job thread while the test reads the results.
*/
class RequestCountingPeerSet : public PeerSet
{
public:
/**
* @param candidates The peers addPeers() may offer, in the order they are
* considered. Fixed at construction. Empty offers no one.
*/
explicit RequestCountingPeerSet(std::vector<std::shared_ptr<Peer>> candidates = {})
: candidates_(std::move(candidates))
{
}
/**
* Offer the candidates to the caller, the way the real peer set offers the
* peers the overlay is tracking.
*
* @param limit The most peers to add, recorded for firstLimit().
* @param hasItem Hard-filters the candidates worth asking, where the real
* peer set only scores with it.
* @param onPeerAdded Called for each selected candidate.
*/
void
addPeers(
std::size_t limit,
std::function<bool(std::shared_ptr<Peer> const&)> hasItem,
std::function<void(std::shared_ptr<Peer> const&)> onPeerAdded) override
{
std::vector<std::shared_ptr<Peer>> selected;
{
std::scoped_lock const lock(mutex_);
if (!firstLimit_)
firstLimit_ = limit;
for (auto const& candidate : candidates_)
{
if (selected.size() >= limit)
break;
// Dedup by tracked id, like the real peer set: onPeerAdded runs once per candidate.
if (hasItem(candidate) && addedPeers_.insert(candidate->id()).second)
selected.push_back(candidate);
}
}
// Outside the lock: onPeerAdded() reenters this object through sendRequest().
for (auto const& peer : selected)
onPeerAdded(peer);
}
/**
* Record the request, and record it in the broadcast count as well when it
* names no peer.
*
* @param peer The peer to ask, or null to ask every tracked peer.
*/
void
sendRequest(
::google::protobuf::Message const&,
protocol::MessageType,
std::shared_ptr<Peer> const& peer) override
{
std::scoped_lock const lock(mutex_);
++requests_;
if (!peer)
++broadcasts_;
}
/**
* The ids of every peer addPeers() has selected, which is what an
* acquisition takes for the peers it is tracking.
*
* Unguarded: every caller of this and of addPeers() is an acquisition
* holding its own mtx_, so the ids stay fixed while a caller iterates. A
* test thread reads addedPeers() instead.
* InboundLedger::getPeerCount() reports zero for these, since it resolves
* ids through the overlay, while these peers live only in this harness.
*
* @return The ids.
*/
[[nodiscard]] std::set<Peer::ID> const&
getPeerIds() const override
{
return addedPeers_;
}
/**
* @return How many requests have been sent through this peer set, counting
* a broadcast as one.
*/
[[nodiscard]] int
requests() const
{
std::scoped_lock const lock(mutex_);
return requests_;
}
/**
* How many of those requests carried no peer of their own.
*
* @return The count.
*/
[[nodiscard]] int
broadcasts() const
{
std::scoped_lock const lock(mutex_);
return broadcasts_;
}
/**
* The limit the first addPeers() call asked for, which is init()'s, since
* onTimer() keeps calling addPeers(1) for as long as an acquisition runs.
*
* @return The limit, or nullopt if addPeers() has not been called.
*/
[[nodiscard]] std::optional<std::size_t>
firstLimit() const
{
std::scoped_lock const lock(mutex_);
return firstLimit_;
}
/**
* A set, because onTimer() keeps re-offering the same candidates.
*
* @return The ids of every peer addPeers() has selected so far.
*/
[[nodiscard]] std::set<Peer::ID>
addedPeers() const
{
std::scoped_lock const lock(mutex_);
return addedPeers_;
}
private:
std::vector<std::shared_ptr<Peer>> const candidates_;
mutable std::mutex mutex_;
int requests_{0};
int broadcasts_{0};
std::optional<std::size_t> firstLimit_;
std::set<Peer::ID> addedPeers_;
};
/**
* The given nodes of a chain as a TMLedgerData, so a test can go through the
* real dispatch. Not in DeepChain, since the protobuf types are xrpld and that
* header is shared with the libxrpl-only gtest binary.
*
* @param chain The chain the nodes came from, which names the reply by default.
* @param data The nodes to include, each with its claimed position.
* @param type The reply type, which selects which map the receiver applies it
* to.
* @param ledgerHash The hash the reply claims to be about, defaulting to the
* chain root for a TX set. A ledger acquisition wants its header hash
* here, since the chain root is only that ledger's account hash.
* @param ledgerSeq The sequence to name in the reply.
* @return The reply packet.
*/
[[nodiscard]] inline std::shared_ptr<protocol::TMLedgerData>
packetFor(
DeepChain const& chain,
std::vector<std::pair<SHAMapNodeID, SHAMapTreeNodePtr>> const& data,
protocol::TMLedgerInfoType type = protocol::liTS_CANDIDATE,
std::optional<UInt256> const& ledgerHash = std::nullopt,
std::uint32_t ledgerSeq = 0)
{
auto packet = std::make_shared<protocol::TMLedgerData>();
auto const hash = ledgerHash.value_or(chain.rootHash.asUInt256());
packet->set_ledgerhash(hash.data(), UInt256::size());
packet->set_ledgerseq(ledgerSeq);
packet->set_type(type);
for (auto const& [nodeID, node] : data)
{
Serializer s;
node->serializeForWire(s);
auto* const ledgerNode = packet->add_nodes();
ledgerNode->set_nodedata(s.peekData().data(), s.peekData().size());
// A leaf carries its own key, so the receiver rebuilds its position from that plus a
// depth. An inner node has no key and needs the full ID. The two fields are a oneof.
if (node->isLeaf())
{
ledgerNode->set_depth(nodeID.getDepth());
}
else
{
ledgerNode->set_id(nodeID.getRawString());
}
}
return packet;
}
/**
* Poll until the condition holds, or give up. An acquisition's own timer and
* the jobs it hands finished work to both run on other threads.
*
* @param condition What to wait for.
* @param deadline The longest to wait.
* @return Whether the condition held before the deadline.
*/
[[nodiscard]] inline bool
waitFor(
std::function<bool()> const& condition,
std::chrono::steady_clock::duration deadline = std::chrono::seconds{10})
{
auto const giveUp = std::chrono::steady_clock::now() + deadline;
while (std::chrono::steady_clock::now() < giveUp)
{
if (condition())
return true;
std::this_thread::sleep_for(std::chrono::milliseconds{10});
}
return condition();
}
/**
* Whether a batch verdict carries exactly the given counts.
*
* The counts, since get() is a log format. It is pinned once, in the
* SHAMapAddNode tests, and is what to pass BEAST_EXPECTS() as the reason a
* check here failed.
*
* @param san The verdict to check.
* @param good How many nodes the batch should have hooked in.
* @param bad How many it should have rejected.
* @param duplicate How many it should have already held.
* @return Whether the verdict matches.
*/
[[nodiscard]] inline bool
tallyIs(SHAMapAddNode const& san, int good, int bad, int duplicate)
{
return san.getGood() == good && san.getBad() == bad && san.getDuplicate() == duplicate;
}
} // namespace xrpl::test

View File

@@ -0,0 +1,537 @@
#include <test/app/AcquireTestHelpers.h>
#include <test/jtx/Env.h>
#include <xrpld/app/ledger/InboundLedger.h>
#include <xrpld/app/ledger/InboundLedgers.h>
#include <xrpld/app/ledger/LedgerMaster.h>
#include <xrpl/basics/base_uint.h>
#include <xrpl/basics/chrono.h>
#include <xrpl/beast/unit_test/suite.h>
#include <xrpl/nodestore/NodeObject.h>
#include <xrpl/protocol/HashPrefix.h>
#include <xrpl/protocol/LedgerHeader.h>
#include <xrpl/protocol/Serializer.h>
#include <xrpl/protocol/jss.h>
#include <xrpl/shamap/SHAMapNodeID.h>
#include <xrpl.pb.h>
#include <chrono>
#include <memory>
#include <mutex>
#include <set>
#include <utility>
#include <vector>
namespace xrpl::test {
/**
* An acquisition that exposes the entry points its bases keep protected, so a
* case can reach them through this subclass alone.
*/
struct TestableInboundLedger final : InboundLedger
{
using InboundLedger::InboundLedger;
/**
* Look for the ledger locally, ask the peers being tracked for the rest,
* and arm the timer, as InboundLedgers::acquire() does. That caller holds
* its collection lock across init(), which releases it, so this stands in
* with a lock of its own. The mutex is declared first, so it outlives the
* lock.
*/
void
startAcquire()
{
std::recursive_mutex collectionMutex;
ScopedLockType collectionLock(collectionMutex);
init(collectionLock);
}
/**
* Ask for more nodes, or judge what has been collected, as a fresh
* acquisition does.
*/
void
triggerAdded()
{
trigger(nullptr, TriggerReason::Added);
}
};
struct InboundLedger_test : public beast::unit_test::Suite
{
/**
* A retry interval short enough that a whole timeout chain costs a fraction
* of a second. TimeoutCounter refuses anything at or below 10ms.
*/
static constexpr auto kFastRetry = std::chrono::milliseconds{20};
/**
* A seed unique to this chain within the suite.
*
* The Env below is shared, and its node store, fetch packs and remembered
* failures are all keyed by hash. A fresh seed per chain gives every chain
* distinct hashes, so each case resolves only its own nodes.
*
* @return The seed.
*/
[[nodiscard]] unsigned int
nextSeed()
{
return ++seed_;
}
/**
* A ledger header naming the given map roots, with a hash derived from the
* fields, so an acquisition accepts the header as its own however the roots
* are chosen.
*
* @param txHash The transaction map root. Zero means no transactions.
* @param accountHash The state map root. Zero is a ledger no acquisition
* can finish.
* @return The header, with its hash filled in.
*/
static LedgerHeader
makeHeader(UInt256 const& txHash, UInt256 const& accountHash)
{
LedgerHeader header;
header.seq = 2;
header.parentCloseTime = NetClock::time_point{};
header.closeTime = NetClock::time_point{};
header.closeTimeResolution = NetClock::duration{10};
header.closeFlags = 0;
header.txHash = txHash;
header.accountHash = accountHash;
header.hash = calculateLedgerHash(header);
return header;
}
/**
* The common shape: no transactions, so only the state map is in play.
*
* @param chain The chain whose root to name as the state hash.
* @return The header, with its hash filled in.
*/
static LedgerHeader
makeHeader(DeepChain const& chain)
{
return makeHeader(UInt256{}, chain.rootHash.asUInt256());
}
/**
* Put the header in the local store, which is the first place tryDB()
* looks. The store keeps its entries, where a fetch pack hands each out
* once, so more than one acquisition of the same ledger can find it.
*
* @param env The environment whose node store to seed.
* @param header The header to store, keyed by its own hash.
*/
static void
storeHeader(jtx::Env& env, LedgerHeader const& header)
{
Serializer s;
s.add32(HashPrefix::LedgerMaster);
addRaw(header, s);
env.app().getNodeFamily().db().store(
NodeObjectType::Ledger, std::move(s.modData()), header.hash, header.seq);
}
/**
* Put every node of a chain in the local store, so a state-map walk
* resolves the whole map from the store.
*
* @param env The environment whose node store to seed.
* @param header The header whose sequence the nodes are stored under.
* @param chain The chain supplying the nodes.
* @param maxDepth The deepest node to store, so a caller can leave a walk
* something to ask for.
*/
static void
storeStateNodes(
jtx::Env& env,
LedgerHeader const& header,
DeepChain const& chain,
unsigned int maxDepth)
{
auto& db = env.app().getNodeFamily().db();
for (auto depth = 0u; depth <= maxDepth; ++depth)
{
db.store(
NodeObjectType::AccountNode,
chain.prefixedNodeAt(depth),
chain.nodeAt(depth)->getHash().asUInt256(),
header.seq);
}
}
/**
* The header alone as a liBASE reply, so takeHeader() alone judges the
* packet.
*
* @param header The header to carry, serialized without a hash prefix, as
* takeHeader() expects.
* @return The reply packet.
*/
static std::shared_ptr<protocol::TMLedgerData>
headerPacket(LedgerHeader const& header)
{
Serializer s;
addRaw(header, s);
auto packet = std::make_shared<protocol::TMLedgerData>();
packet->set_ledgerhash(header.hash.data(), UInt256::size());
packet->set_ledgerseq(header.seq);
packet->set_type(protocol::liBASE);
packet->add_nodes()->set_nodedata(s.peekData().data(), s.peekData().size());
return packet;
}
/**
* An acquisition started with nothing in the store for its hash, so it
* waits on its peers.
*
* @param env The environment to run in.
* @param header The header of the ledger to acquire.
* @return The acquisition, started.
*/
static std::shared_ptr<TestableInboundLedger>
startPeerAcquire(jtx::Env& env, LedgerHeader const& header)
{
auto acquire = std::make_shared<TestableInboundLedger>(
env.app(),
header.hash,
header.seq,
InboundLedger::Reason::GENERIC,
stopwatch(),
std::make_unique<RequestCountingPeerSet>());
acquire->startAcquire();
return acquire;
}
/**
* A ledger whose maps all resolve locally finishes on the spot, and the
* finished ledger is immutable and handed on.
*
* Both entry points into tryDB() are driven: InboundLedgers::acquire(), the
* only caller of init(), and checkLocal(), which reaches done().
*
* @param env The environment to run in.
*/
void
testLocalLedgerCompletesAcquire(jtx::Env& env)
{
testcase("A ledger found locally completes the acquire");
// A chain ending in a real leaf, so the state map is complete rather than merely rooted.
// No transactions, so only the state map is in play.
auto const chain = DeepChain::toLeaf(2, nextSeed());
auto const header = makeHeader(chain);
storeHeader(env, header);
storeStateNodes(env, header, chain, chain.deepestDepth);
// acquire() returns the ledger only once the acquisition is complete and unfailed, so a
// non-null result shows tryDB() found it in the store.
auto const acquired = env.app().getInboundLedgers().acquire(
header.hash, header.seq, InboundLedger::Reason::GENERIC);
BEAST_EXPECT(acquired != nullptr);
if (acquired)
{
BEAST_EXPECT(acquired->isImmutable());
BEAST_EXPECT(acquired->header().hash == header.hash);
}
// init() hands a ledger it completed to LedgerMaster itself, which is what makes it
// available to everything else.
BEAST_EXPECT(env.app().getLedgerMaster().getLedgerByHash(header.hash) != nullptr);
// The failure list stays clear, which is the other arm of done().
BEAST_EXPECT(!env.app().getInboundLedgers().isFailure(header.hash));
// The same ledger through checkLocal(), which unlike init() reaches done(). Everything it
// needs is still in the store, since the first acquisition read rather than consumed it.
auto again = std::make_shared<InboundLedger>(
env.app(),
header.hash,
header.seq,
InboundLedger::Reason::GENERIC,
stopwatch(),
std::make_unique<RequestCountingPeerSet>());
// True because the acquisition ended, which it reports only after done() has run.
BEAST_EXPECT(again->checkLocal());
BEAST_EXPECT(again->isComplete());
BEAST_EXPECT(!again->isFailed());
auto const settled = again->getLedger();
BEAST_EXPECT(settled != nullptr);
if (settled)
BEAST_EXPECT(settled->isImmutable());
BEAST_EXPECT(!env.app().getInboundLedgers().isFailure(header.hash));
}
/**
* An acquisition that fails on local data must still signal.
*
* Both entry points that reach tryDB() are covered. done() is what
* signals, runs logFailure(), and lands the hash in recentFailures_,
* which stops the next round asking for the same ledger.
* recentFailures_ is what the assertions watch, since it is the
* caller-visible consequence of having signaled.
*
* @param env The environment to run in.
*/
void
testLocalFailureSignalsDone(jtx::Env& env)
{
testcase("An acquisition that fails locally still signals");
// A zero account hash is a ledger no acquisition can ever finish, and tryDB() says so as
// soon as it has the header.
auto const header = makeHeader(UInt256{}, UInt256{});
storeHeader(env, header);
BEAST_EXPECT(!env.app().getInboundLedgers().isFailure(header.hash));
// acquire() is the only caller of init(), and returns nullptr for a failed acquisition.
BEAST_EXPECT(
env.app().getInboundLedgers().acquire(
header.hash, header.seq, InboundLedger::Reason::GENERIC) == nullptr);
// The failure reached recentFailures_, which is what stops the next round asking again.
BEAST_EXPECT(waitFor([&] { return env.app().getInboundLedgers().isFailure(header.hash); }));
// The other way tryDB() judges a ledger unobtainable: the header it found is not the one
// asked for. A non-zero account hash, so the route above cannot be what fails this one.
auto const strayHeader = makeHeader(UInt256{}, UInt256{2});
storeHeader(env, strayHeader);
BEAST_EXPECT(!env.app().getInboundLedgers().isFailure(strayHeader.hash));
// acquire() hands the sequence to the acquisition unscreened, so one past the stored
// header's is enough to make tryDB() reject what it reads.
auto const wrongSeq = strayHeader.seq + 1;
BEAST_EXPECT(
env.app().getInboundLedgers().acquire(
strayHeader.hash, wrongSeq, InboundLedger::Reason::GENERIC) == nullptr);
// tryDB() discards the ledger it built before it gives up, so done() runs here with none
// to report. The failure is recorded all the same, which is what the call is for.
BEAST_EXPECT(
waitFor([&] { return env.app().getInboundLedgers().isFailure(strayHeader.hash); }));
// The other route into tryDB(): a trigger() on an acquisition that has no header yet. A
// hash of its own, so this drives a fresh entry.
auto const otherHeader = makeHeader(UInt256{1}, UInt256{});
storeHeader(env, otherHeader);
auto viaTrigger = std::make_shared<TestableInboundLedger>(
env.app(),
otherHeader.hash,
otherHeader.seq,
InboundLedger::Reason::GENERIC,
stopwatch(),
std::make_unique<RequestCountingPeerSet>());
BEAST_EXPECT(!env.app().getInboundLedgers().isFailure(otherHeader.hash));
viaTrigger->triggerAdded();
// Failed before any header state was published: no ledger, and have_header false.
BEAST_EXPECT(viaTrigger->isFailed());
BEAST_EXPECT(!viaTrigger->isComplete());
BEAST_EXPECT(viaTrigger->getLedger() == nullptr);
BEAST_EXPECT(!viaTrigger->getJson(0)[jss::have_header].asBool());
BEAST_EXPECT(
waitFor([&] { return env.app().getInboundLedgers().isFailure(otherHeader.hash); }));
}
/**
* A peer-supplied header whose account hash is zero fails the acquisition
* and signals, as a local one does. The peer is not charged, and
* recentFailures_ is what shows done() ran.
*
* @param env The environment to run in.
*/
void
testPeerZeroAccountHashFails(jtx::Env& env)
{
testcase("A peer header with a zero account hash fails the acquire");
// A real transaction root and a zero state root. Nothing is stored for this hash, so only a
// peer can supply the header.
auto const chain = DeepChain::toLeaf(2, nextSeed());
auto const header = makeHeader(chain.rootHash.asUInt256(), UInt256{});
BEAST_EXPECT(!env.app().getInboundLedgers().isFailure(header.hash));
auto acquire = startPeerAcquire(env, header);
BEAST_EXPECT(!acquire->isFailed());
BEAST_EXPECT(!acquire->getJson(0)[jss::have_header].asBool());
auto const peer = std::make_shared<ChargeRecordingPeer>();
BEAST_EXPECT(acquire->gotData(peer, headerPacket(header)));
acquire->runData();
// Failed, not complete, and the header is dropped, as tryDB() leaves one.
BEAST_EXPECT(acquire->isFailed());
BEAST_EXPECT(!acquire->isComplete());
BEAST_EXPECT(acquire->getLedger() == nullptr);
BEAST_EXPECT(!acquire->getJson(0)[jss::have_header].asBool());
// The header is the one asked for, so the sender is not to blame.
BEAST_EXPECT(peer->charges().empty());
// done() ran, so the hash reached recentFailures_, which stops the next round asking again.
BEAST_EXPECT(waitFor([&] { return env.app().getInboundLedgers().isFailure(header.hash); }));
}
/**
* A peer-supplied header with a zero transaction hash is taken, since that
* is an empty transaction set, and the acquisition finishes once the peer
* supplies the state map.
*
* @param env The environment to run in.
*/
void
testPeerHeaderWithoutTransactionsCompletes(jtx::Env& env)
{
testcase("A peer header with no transactions completes the acquire");
// No transactions, and a state map that ends in a real leaf, so the acquisition can finish.
auto const chain = DeepChain::toLeaf(2, nextSeed());
auto const header = makeHeader(chain);
auto acquire = startPeerAcquire(env, header);
auto const peer = std::make_shared<ChargeRecordingPeer>();
BEAST_EXPECT(acquire->gotData(peer, headerPacket(header)));
acquire->runData();
// The header is taken, the empty transaction map counts as fetched, and the state map is
// what remains.
BEAST_EXPECT(!acquire->isFailed());
BEAST_EXPECT(!acquire->isComplete());
auto const json = acquire->getJson(0);
BEAST_EXPECT(json[jss::have_header].asBool());
BEAST_EXPECT(json[jss::have_transactions].asBool());
BEAST_EXPECT(!json[jss::have_state].asBool());
// The whole state map as state-node replies, root first, which completes the acquisition.
// The second gotData() answers false, since the first already asked for a dispatch.
BEAST_EXPECT(acquire->gotData(
peer,
packetFor(
chain,
{{SHAMapNodeID{}, chain.nodeAt(0)}},
protocol::liAS_NODE,
header.hash,
header.seq)));
BEAST_EXPECT(!acquire->gotData(
peer,
packetFor(
chain, chain.nodesBelowRoot(), protocol::liAS_NODE, header.hash, header.seq)));
acquire->runData();
BEAST_EXPECT(acquire->isComplete());
BEAST_EXPECT(!acquire->isFailed());
BEAST_EXPECT(peer->charges().empty());
auto const settled = acquire->getLedger();
BEAST_EXPECT(settled != nullptr);
if (settled)
BEAST_EXPECT(settled->isImmutable());
BEAST_EXPECT(!env.app().getInboundLedgers().isFailure(header.hash));
}
/**
* The retry timer re-asks, then gives up and signals.
*
* The only case that drives onTimer() rather than trigger() directly,
* which is what covers the give-up: past kLedgerTimeoutRetriesMax the
* acquisition fails itself and done() records that, so the same
* doomed ledger is not asked for again on the next round. It is also
* what the retry interval is a constructor parameter for, since the
* chain runs past kLedgerTimeoutRetriesMax ticks of three seconds
* apiece in production.
*
* The store is empty for this hash and no reply arrives, so every tick
* counts a timeout and the count climbs to the limit. A hash of its own,
* so this case owns its entry in the failure list.
*
* @param env The environment to run in.
*/
void
testTimerRetriesThenGivesUp(jtx::Env& env)
{
testcase("The retry timer re-asks, then gives up");
UInt256 const kUnknownLedger{8};
// One candidate, which onTimer() re-offers on every tick.
auto const candidate = std::make_shared<ChargeRecordingPeer>();
auto peerSet =
std::make_unique<RequestCountingPeerSet>(std::vector<std::shared_ptr<Peer>>{candidate});
auto* const peerSetPtr = peerSet.get();
auto acquire = std::make_shared<TestableInboundLedger>(
env.app(),
kUnknownLedger,
0,
InboundLedger::Reason::GENERIC,
stopwatch(),
std::move(peerSet),
kFastRetry);
BEAST_EXPECT(!env.app().getInboundLedgers().isFailure(kUnknownLedger));
// init() finds the hash absent from the store, so it asks the candidate and arms the retry
// timer. The case counts from here, past those first requests.
acquire->startAcquire();
BEAST_EXPECT(!acquire->isFailed());
int const requestsFromInit = peerSetPtr->requests();
BEAST_EXPECT(requestsFromInit > 0);
BEAST_EXPECT(peerSetPtr->addedPeers() == std::set<Peer::ID>{candidate->id()});
// Every tick asks again, and past kLedgerTimeoutRetriesMax (6) the chain gives up.
BEAST_EXPECT(waitFor([&] { return acquire->isFailed(); }));
BEAST_EXPECT(!acquire->isComplete());
BEAST_EXPECT(peerSetPtr->requests() > requestsFromInit);
// done() remembered the hash, which is what stops the next round asking again.
BEAST_EXPECT(
waitFor([&] { return env.app().getInboundLedgers().isFailure(kUnknownLedger); }));
}
void
run() override
{
// One Env for the suite, since building one costs far more than any case here. Safe
// because every chain is seeded through nextSeed(): the node store, the fetch packs and
// the remembered failures are all shared, and all three are keyed by hash.
jtx::Env env{*this};
testLocalLedgerCompletesAcquire(env);
testLocalFailureSignalsDone(env);
testPeerZeroAccountHashFails(env);
testPeerHeaderWithoutTransactionsCompletes(env);
// Last: the only case that waits out a whole timeout chain.
testTimerRetriesThenGivesUp(env);
}
private:
unsigned int seed_{0};
};
BEAST_DEFINE_TESTSUITE(InboundLedger, app, xrpl);
} // namespace xrpl::test

View File

@@ -0,0 +1,490 @@
#include <test/app/AcquireTestHelpers.h>
#include <test/jtx/Env.h>
#include <xrpld/app/ledger/InboundTransactions.h>
#include <xrpld/app/ledger/detail/TransactionAcquire.h>
#include <xrpld/overlay/Peer.h>
#include <xrpl/basics/base_uint.h>
#include <xrpl/beast/unit_test/suite.h>
#include <xrpl/resource/Fees.h>
#include <xrpl/shamap/SHAMap.h>
#include <xrpl/shamap/SHAMapAddNode.h>
#include <xrpl/shamap/SHAMapNodeID.h>
#include <xrpl.pb.h>
#include <chrono>
#include <memory>
#include <set>
#include <utility>
#include <vector>
namespace xrpl::test {
/**
* An acquisition that exposes state its bases keep protected.
*/
struct TestableTransactionAcquire final : TransactionAcquire
{
using TransactionAcquire::TransactionAcquire;
/**
* Whether the set being acquired is still structurally coherent.
*
* @return Whether the map is still valid.
*/
[[nodiscard]] bool
isMapValid() const
{
// Under the lock: a batch on another thread can reach the verdict.
ScopedLockType const sl(mtx_);
return map_->isValid();
}
/**
* Whether a batch has advanced the set since the flag was last cleared.
*
* @return Whether progress has been recorded.
*/
[[nodiscard]] bool
madeProgress() const
{
// Under the lock: a timer tick clears the flag on a job thread.
ScopedLockType const sl(mtx_);
return progress_;
}
/**
* Forget any recorded progress.
*/
void
clearProgress()
{
ScopedLockType const sl(mtx_);
progress_ = false;
}
};
struct TransactionAcquire_test : public beast::unit_test::Suite
{
/**
* A retry interval short enough that a whole timeout chain costs a fraction
* of a second.
*
* TimeoutCounter refuses anything at or below 10ms. At this interval the
* window between the first retry (four timeouts in) and giving up (twenty)
* is still a third of a second, which is what the one case that watches
* both needs.
*/
static constexpr auto kFastRetry = std::chrono::milliseconds{20};
/**
* A seed no other chain in this suite has used.
*
* The Env below is shared, and ConsensusTransSetSF::gotNode() puts
* every node it accepts into the application-wide NodeCache while
* InboundTransactions keys its acquisitions by set hash. A fresh seed
* per chain gives every chain distinct hashes, so each case resolves
* and revives only its own.
*
* @return The seed.
*/
[[nodiscard]] unsigned int
nextSeed()
{
return ++seed_;
}
/**
* Whether takeNodes() declined to look at the data at all.
*
* TimeoutCounter::complete_ and failed_ are both protected, so this
* stands in for either: a done acquisition returns a verdict with every
* count at zero.
*
* @param san The verdict a takeNodes() call returned.
* @return Whether that verdict shows the data was never looked at.
*/
static bool
wasIgnored(SHAMapAddNode const& san)
{
return tallyIs(san, 0, 0, 0);
}
/**
* Wait for a finished set to reach InboundTransactions.
*
* done() hands the map over through a job, so a delivered set is what
* separates completion from a stop.
*
* @param env The environment whose InboundTransactions to watch.
* @param setHash The set to wait for.
* @return The delivered map, or nullptr if none arrived.
*/
[[nodiscard]] static std::shared_ptr<SHAMap>
waitForDeliveredSet(jtx::Env& env, UInt256 const& setHash)
{
auto& inbound = env.app().getInboundTransactions();
// acquire=false keeps this to a lookup rather than registering an acquisition. The job is
// queued before takeNodes() returns, so a few seconds is a generous deadline.
std::shared_ptr<SHAMap> delivered;
if (!waitFor(
[&] { return (delivered = inbound.getSet(setHash, false)) != nullptr; },
std::chrono::seconds{5}))
return nullptr;
return delivered;
}
/**
* A chain ending in a leaf completes the acquisition, which stops the
* asking and hands the map to InboundTransactions. Later replies for it
* are then left alone.
*
* @param env The environment to run in.
*/
void
testHappyPathCompletesAcquisition(jtx::Env& env)
{
testcase("A chain ending in a leaf completes the acquire");
auto const chain = DeepChain::toLeaf(3, nextSeed());
auto peerSet = std::make_unique<RequestCountingPeerSet>();
auto* const peerSetPtr = peerSet.get();
UInt256 const setHash = chain.rootHash.asUInt256();
auto const acquire =
std::make_shared<TransactionAcquire>(env.app(), setHash, std::move(peerSet));
auto const peer = std::make_shared<ChargeRecordingPeer>();
// The root alone leaves the set incomplete, so accepting it asks for the level below.
auto const rootResult = acquire->takeNodes({{SHAMapNodeID{}, chain.nodeAt(0)}}, peer);
BEAST_EXPECT(rootResult.isUseful());
int const requestsWhileIncomplete = peerSetPtr->requests();
BEAST_EXPECT(requestsWhileIncomplete > 0);
// The rest of the chain, ending in the leaf.
auto const result = acquire->takeNodes(chain.nodesBelowRoot(), peer);
BEAST_EXPECT(result.isUseful());
BEAST_EXPECT(!result.isInvalid());
// The request count holds steady, so the asking has stopped.
BEAST_EXPECT(peerSetPtr->requests() == requestsWhileIncomplete);
auto const delivered = waitForDeliveredSet(env, setHash);
BEAST_EXPECT(delivered != nullptr);
if (delivered)
{
BEAST_EXPECT(delivered->getHash() == chain.rootHash);
BEAST_EXPECT(delivered->isValid());
}
// A reply arriving after the set is finished is left alone, and the asking stays stopped.
// init() stayed uncalled, so the set above is what settled the acquisition.
BEAST_EXPECT(wasIgnored(acquire->takeNodes({{SHAMapNodeID{}, chain.nodeAt(0)}}, peer)));
BEAST_EXPECT(peerSetPtr->requests() == requestsWhileIncomplete);
}
/**
* Two peers each answering with a different missing piece are both accepted
* without penalty, and the set completes from their combined replies.
*
* Driven through InboundTransactions::gotData(), so the leaf goes over the
* real dispatch, which derives a leaf's position from its own key.
*
* @param env The environment to run in.
*/
void
testTwoPeersEachSupplyPartOfTheSet(jtx::Env& env)
{
testcase("Two peers each supplying part of a set are both accepted without penalty");
auto const chain = DeepChain::toLeaf(3, nextSeed());
auto& inbound = env.app().getInboundTransactions();
// acquire=true registers the TransactionAcquire that gotData() looks up by hash.
UInt256 const setHash = chain.rootHash.asUInt256();
BEAST_EXPECT(inbound.getSet(setHash, true) == nullptr);
// The first peer answers with the root only.
auto const peerA = std::make_shared<ChargeRecordingPeer>();
inbound.gotData(setHash, peerA, packetFor(chain, {{SHAMapNodeID{}, chain.nodeAt(0)}}));
BEAST_EXPECT(peerA->charges().empty());
// The second answers with everything the first left out, and finishes the set.
auto const peerB = std::make_shared<ChargeRecordingPeer>();
inbound.gotData(setHash, peerB, packetFor(chain, chain.nodesBelowRoot()));
BEAST_EXPECT(peerB->charges().empty());
auto const delivered = waitForDeliveredSet(env, setHash);
BEAST_EXPECT(delivered != nullptr);
if (delivered)
BEAST_EXPECT(delivered->getHash() == chain.rootHash);
}
/**
* A root that does not hash to the set we asked for is a plain mismatch,
* and has to leave the acquisition able to try another peer.
*
* @param env The environment to run in.
*/
void
testBadRootKeepsAcquireAlive(jtx::Env& env)
{
testcase("A mismatched root leaves the acquire recoverable");
DeepChain const chain{nextSeed()};
// Acquire an unrelated hash, so the chain's root cannot match it.
auto const acquire = std::make_shared<TestableTransactionAcquire>(
env.app(), UInt256{42}, std::make_unique<RequestCountingPeerSet>());
auto const peer = std::make_shared<ChargeRecordingPeer>();
auto const result = acquire->takeNodes({{SHAMapNodeID{}, chain.nodeAt(0)}}, peer);
BEAST_EXPECTS(tallyIs(result, 0, 1, 0), result.get());
// A mismatched root tells the acquisition only that this peer's answer is wrong, so the
// map keeps its state.
BEAST_EXPECT(acquire->isMapValid());
// Still alive: the next packet is examined rather than waved through.
BEAST_EXPECT(!wasIgnored(acquire->takeNodes({{SHAMapNodeID{}, chain.nodeAt(0)}}, peer)));
}
/**
* A second reply carrying a root we already have must stay free.
*
* This is what an honest second responder to the initial fan-out sends:
* trigger() broadcasts to every tracked peer, so several answer the same
* request and all but the first carry data the map already holds, so they
* go uncharged.
*
* Covers the root specifically, which takeNodes() short-circuits on
* haveRoot_ without consulting the map. A repeated non-root node takes the
* other route, through addKnownNode() - see
* testDuplicateNonRootReplyIsFree().
*
* @param env The environment to run in.
*/
void
testDuplicateRootReplyIsFree(jtx::Env& env)
{
testcase("A reply of a root we already have is free");
DeepChain const chain{nextSeed()};
auto& inbound = env.app().getInboundTransactions();
// acquire=true registers the TransactionAcquire that gotData() looks up by hash.
UInt256 const setHash = chain.rootHash.asUInt256();
BEAST_EXPECT(inbound.getSet(setHash, true) == nullptr);
auto const rootPacket = packetFor(chain, {{SHAMapNodeID{}, chain.nodeAt(0)}});
// The first responder supplies the root, which is genuinely useful.
auto const firstPeer = std::make_shared<ChargeRecordingPeer>();
inbound.gotData(setHash, firstPeer, rootPacket);
BEAST_EXPECT(firstPeer->charges().empty());
// The second sends the same root. It answered what we asked, so it goes uncharged.
auto const secondPeer = std::make_shared<ChargeRecordingPeer>();
inbound.gotData(setHash, secondPeer, rootPacket);
BEAST_EXPECT(secondPeer->charges().empty());
}
/**
* A repeated non-root node stays free too.
*
* The counterpart to testDuplicateRootReplyIsFree(), covering the route
* that consults the map: addKnownNode() reports a node it already holds as
* a duplicate, which isGood() counts as success.
*
* @param env The environment to run in.
*/
void
testDuplicateNonRootReplyIsFree(jtx::Env& env)
{
testcase("A repeated non-root node is free");
auto const chain = DeepChain::toLeaf(3, nextSeed());
auto& inbound = env.app().getInboundTransactions();
UInt256 const setHash = chain.rootHash.asUInt256();
BEAST_EXPECT(inbound.getSet(setHash, true) == nullptr);
auto const rootPeer = std::make_shared<ChargeRecordingPeer>();
inbound.gotData(setHash, rootPeer, packetFor(chain, {{SHAMapNodeID{}, chain.nodeAt(0)}}));
BEAST_EXPECT(rootPeer->charges().empty());
// Depth 1 alone, so the set stays incomplete and the acquisition keeps examining data
// rather than waving the second copy through as a late reply.
auto const level1 = packetFor(chain, {{chain.idAt(1), chain.nodeAt(1)}});
auto const firstPeer = std::make_shared<ChargeRecordingPeer>();
inbound.gotData(setHash, firstPeer, level1);
BEAST_EXPECT(firstPeer->charges().empty());
auto const secondPeer = std::make_shared<ChargeRecordingPeer>();
inbound.gotData(setHash, secondPeer, level1);
BEAST_EXPECT(secondPeer->charges().empty());
}
/**
* A reply whose node data cannot be deserialized is charged for.
*
* gotData() rejects the packet before the acquisition is handed
* anything, so this pins the charge on the dispatch layer rather than
* on takeNodes(). It also gives this suite's "was not charged"
* assertions their teeth: a harness that recorded no charge at all
* would satisfy all of them and fail only here.
*
* @param env The environment to run in.
*/
void
testUndeserializableNodeIsCharged(jtx::Env& env)
{
testcase("A reply with undeserializable node data is charged");
DeepChain const chain{nextSeed()};
auto& inbound = env.app().getInboundTransactions();
UInt256 const setHash = chain.rootHash.asUInt256();
BEAST_EXPECT(inbound.getSet(setHash, true) == nullptr);
// getTreeNode() rejects this before the acquisition is handed anything.
auto packet = std::make_shared<protocol::TMLedgerData>();
packet->set_ledgerhash(setHash.data(), UInt256::size());
packet->set_ledgerseq(0);
packet->set_type(protocol::liTS_CANDIDATE);
auto* const node = packet->add_nodes();
node->set_nodedata("\xff", 1);
node->set_id(SHAMapNodeID{}.getRawString());
auto const peer = std::make_shared<ChargeRecordingPeer>();
inbound.gotData(setHash, peer, packet);
BEAST_EXPECT(peer->charges() == std::vector{resource::kFeeInvalidData});
}
/**
* init() passes hasTxSet() to addPeers() as its candidate filter.
*
* init() hands addPeers() hasTxSet(hash_) as its hasItem callback and
* trigger() as its onPeerAdded callback. The harness selects on that
* callback while the real peer set only scores with it, so this pins which
* callbacks the acquisition supplies, not how many peers production asks.
*
* @param env The environment to run in.
*/
void
testInitFiltersCandidatesByHasTxSet(jtx::Env& env)
{
testcase("init() passes hasTxSet as its candidate filter");
DeepChain const chain{nextSeed()};
// Ordered with the useless peer first, so a filter that is ignored altogether shows up as
// the wrong peer being asked rather than as one extra request.
auto const withoutSet = std::make_shared<ChargeRecordingPeer>(false);
auto const withSet = std::make_shared<ChargeRecordingPeer>(true);
auto peerSet = std::make_unique<RequestCountingPeerSet>(
std::vector<std::shared_ptr<Peer>>{withoutSet, withSet});
auto* const peerSetPtr = peerSet.get();
auto const acquire = std::make_shared<TransactionAcquire>(
env.app(), chain.rootHash.asUInt256(), std::move(peerSet));
static constexpr int kStartPeers = 2;
acquire->init(kStartPeers);
// Stop the retry loop, which keeps offering the same candidates while the acquisition runs.
acquire->cancel();
BEAST_EXPECT(peerSetPtr->firstLimit() == kStartPeers);
BEAST_EXPECT(peerSetPtr->addedPeers() == std::set<Peer::ID>{withSet->id()});
// The peer that was added is also asked, rather than merely tracked.
BEAST_EXPECT(peerSetPtr->requests() >= 1);
}
/**
* The retry timer broadcasts with no peer of its own, then gives up on
* its own.
*
* The only case that reaches onTimer(). Pins the two behaviors, not the
* thresholds they trip at: bounding those means asserting on wall clock.
* Both are read in one poll, so a fast interval cannot let the give-up land
* between the two readings.
*
* The acquisition tracks no peer, so the broadcast this case counts reaches
* nobody in production either. What is pinned is that onTimer() issues it.
*
* @param env The environment to run in.
*/
void
testTimerBroadcastsThenGivesUp(jtx::Env& env)
{
testcase("The retry timer broadcasts, then gives up");
DeepChain const chain{nextSeed()};
auto peerSet = std::make_unique<RequestCountingPeerSet>();
auto* const peerSetPtr = peerSet.get();
// An unrelated hash, so the probe below is always rejected.
auto const acquire = std::make_shared<TransactionAcquire>(
env.app(), UInt256{42}, std::move(peerSet), kFastRetry);
// No candidates, so nothing goes out until onTimer() decides to broadcast.
acquire->init(1);
BEAST_EXPECT(peerSetPtr->requests() == 0);
// A root that cannot hash to this acquisition's set is rejected without recording progress,
// so polling with it does not postpone the timeout. A fresh peer each time keeps the
// rejections from piling up on one.
auto const probe = [&] {
return wasIgnored(acquire->takeNodes(
{{SHAMapNodeID{}, chain.nodeAt(0)}}, std::make_shared<ChargeRecordingPeer>()));
};
// kNormTimeouts (4) intervals in, onTimer() broadcasts with no peer of its own to ask, and
// the acquisition is still examining data. Read as a broadcast rather than as a request,
// since a request here would mean a peer was named.
BEAST_EXPECT(waitFor([&] { return peerSetPtr->broadcasts() > 0 && !probe(); }));
// Past kMaxTimeouts (20) it fails itself, and stops examining data.
BEAST_EXPECT(waitFor(probe));
// The count outlives the poll that saw the broadcast.
BEAST_EXPECT(peerSetPtr->broadcasts() > 0);
}
void
run() override
{
// One Env for the suite, which is safe only because every chain is seeded through
// nextSeed(): cases sharing an Env share a NodeCache, so they must not share a hash.
jtx::Env env{*this};
testHappyPathCompletesAcquisition(env);
testTwoPeersEachSupplyPartOfTheSet(env);
testBadRootKeepsAcquireAlive(env);
testDuplicateRootReplyIsFree(env);
testDuplicateNonRootReplyIsFree(env);
testUndeserializableNodeIsCharged(env);
testInitFiltersCandidatesByHasTxSet(env);
// Last: the only case that waits out a whole timeout chain.
testTimerBroadcastsThenGivesUp(env);
}
private:
unsigned int seed_{0};
};
BEAST_DEFINE_TESTSUITE(TransactionAcquire, app, xrpl);
} // namespace xrpl::test

View File

@@ -1,123 +0,0 @@
#pragma once
#include <xrpl/basics/ByteUtilities.h>
#include <xrpl/basics/base_uint.h>
#include <xrpl/basics/chrono.h>
#include <xrpl/basics/contract.h>
#include <xrpl/beast/utility/Journal.h>
#include <xrpl/config/BasicConfig.h>
#include <xrpl/config/Constants.h>
#include <xrpl/nodestore/Database.h>
#include <xrpl/nodestore/DummyScheduler.h>
#include <xrpl/nodestore/Manager.h>
#include <xrpl/shamap/Family.h>
#include <xrpl/shamap/FullBelowCache.h>
#include <xrpl/shamap/TreeNodeCache.h>
#include <cstdint>
#include <memory>
#include <stdexcept>
namespace xrpl::test {
/**
* Test implementation of Family for unit tests.
*
* Uses an in-memory NodeStore database and simple caches.
* The missingNode methods throw since tests shouldn't encounter missing nodes.
*/
class TestFamily : public Family
{
private:
std::unique_ptr<node_store::Database> db_;
TestStopwatch clock_;
std::shared_ptr<FullBelowCache> fbCache_;
std::shared_ptr<TreeNodeCache> tnCache_;
node_store::DummyScheduler scheduler_;
beast::Journal j_;
public:
explicit TestFamily(beast::Journal j)
: fbCache_(std::make_shared<FullBelowCache>("TestFamily full below cache", clock_, j))
, tnCache_(
std::make_shared<TreeNodeCache>(
"TestFamily tree node cache",
65536,
std::chrono::minutes{1},
clock_,
j))
, j_(j)
{
Section config;
config.set(Keys::kType, "memory");
config.set(Keys::kPath, "TestFamily");
db_ = node_store::Manager::instance().makeDatabase(megabytes(4), scheduler_, 1, config, j);
}
node_store::Database&
db() override
{
return *db_;
}
[[nodiscard]] node_store::Database const&
db() const override
{
return *db_;
}
beast::Journal const&
journal() override
{
return j_;
}
std::shared_ptr<FullBelowCache>
getFullBelowCache() override
{
return fbCache_;
}
std::shared_ptr<TreeNodeCache>
getTreeNodeCache() override
{
return tnCache_;
}
void
sweep() override
{
fbCache_->sweep();
tnCache_->sweep();
}
void
missingNodeAcquireBySeq(std::uint32_t refNum, UInt256 const& nodeHash) override
{
Throw<std::runtime_error>("TestFamily: missing node (by seq)");
}
void
missingNodeAcquireByHash(UInt256 const& refHash, std::uint32_t refNum) override
{
Throw<std::runtime_error>("TestFamily: missing node (by hash)");
}
void
reset() override
{
(*fbCache_).reset();
(*tnCache_).reset();
}
/**
* Access the test clock for time manipulation in tests.
*/
TestStopwatch&
clock()
{
return clock_;
}
};
} // namespace xrpl::test

View File

@@ -12,8 +12,8 @@
#include <boost/asio/io_context.hpp>
#include <helpers/TestFamily.h>
#include <helpers/TestSink.h>
#include <shamap/common.h>
#include <cstdint>
#include <memory>
@@ -73,7 +73,7 @@ class TestServiceRegistry : public ServiceRegistry
{
TestLogs logs_{beast::Severity::Warning};
boost::asio::io_context ioContext_;
TestFamily family_{logs_.journal("TestFamily")};
tests::TestNodeFamily family_{logs_.journal("TestNodeFamily")};
LoadFeeTrack feeTrack_{logs_.journal("LoadFeeTrack")};
TestNetworkIDService networkIDService_;
HashRouter hashRouter_{HashRouter::Setup{}, stopwatch()};

View File

@@ -0,0 +1,268 @@
#pragma once
#include <xrpl/basics/Blob.h>
#include <xrpl/basics/SHAMapHash.h>
#include <xrpl/basics/Slice.h>
#include <xrpl/basics/base_uint.h>
#include <xrpl/basics/contract.h>
#include <xrpl/protocol/Serializer.h>
#include <xrpl/shamap/SHAMap.h>
#include <xrpl/shamap/SHAMapAddNode.h>
#include <xrpl/shamap/SHAMapLeafNode.h>
#include <xrpl/shamap/SHAMapNodeID.h>
#include <xrpl/shamap/SHAMapTreeNode.h>
#include <cstddef>
#include <optional>
#include <stdexcept>
#include <utility>
#include <vector>
namespace xrpl::tests {
/**
* A chain of inner nodes in wire form, from the root down to one deepest node,
* each with one real child.
*
* Built bottom-up, so every node hashes correctly and the root hash commits to
* the whole shape. Each node sits on the branch pathKey selects at its depth.
*
* Two shapes: the constructors run inner nodes all the way to
* SHAMap::kLeafDepth, a depth only a leaf may occupy, and toLeaf() ends at a
* real transaction leaf.
*
* Depends only on libxrpl, so either test tree can include it. The xrpld
* counterpart is test::packetFor() in src/test/app/AcquireTestHelpers.h.
*/
struct DeepChain
{
// nodes[d] is the deserialized node for depth d.
std::vector<SHAMapTreeNodePtr> nodes;
SHAMapHash rootHash;
// The key whose path through the tree this chain spells out. Zero for a chain
// built without a leaf, which therefore sits on branch 0 at every depth.
UInt256 pathKey;
// The depth of the deepest node, the last one nodesBelowRoot() hands out.
unsigned int deepestDepth{SHAMap::kLeafDepth};
/**
* The payload size of the leaf toLeaf() builds, which is the smallest a
* SHAMap item may be.
*/
static constexpr std::size_t kLeafItemBytes = kMinShaMapItemBytes;
/**
* A chain of inner nodes reaching SHAMap::kLeafDepth, a depth only a leaf
* may occupy.
*
* @param seed Varies the whole chain, so two chains coexist with distinct
* nodes. Caches and fetch packs are keyed by hash, so
* identically-seeded chains are the same chain.
*/
explicit DeepChain(unsigned int seed = 1) : DeepChain(std::nullopt, seed)
{
}
/**
* A chain ending in a real transaction leaf, which completes an
* acquisition.
*
* @param depth Where the leaf sits, at most SHAMap::kLeafDepth. Zero puts
* the leaf at the root. A deeper value throws std::logic_error, since
* no leaf can sit there.
* @param seed Varies the leaf's contents, and so the whole chain. See the
* constructor.
* @return The chain.
*/
[[nodiscard]] static DeepChain
toLeaf(unsigned int depth, unsigned int seed = 1)
{
if (depth > SHAMap::kLeafDepth)
Throw<std::logic_error>("DeepChain: leaf depth past SHAMap::kLeafDepth");
return DeepChain{std::optional{depth}, seed};
}
/**
* The node the chain holds at the given depth, root first.
*
* @param depth The depth of the node to return, at most deepestDepth.
* @return The node.
*/
[[nodiscard]] SHAMapTreeNodePtr
nodeAt(unsigned int depth) const
{
return nodes[depth];
}
/**
* Where the node at the given depth claims to belong, which is on the
* path to pathKey.
*
* @param depth The depth of the node to locate.
* @return The node's claimed position.
*/
[[nodiscard]] SHAMapNodeID
idAt(unsigned int depth) const
{
return SHAMapNodeID::createID(depth, pathKey);
}
/**
* The same node in the prefixed form used for storage and fetch packs,
* which is what hashes to the node's own hash.
*
* @param depth The depth of the node to serialize.
* @return The node's prefixed serialized form.
*/
[[nodiscard]] Blob
prefixedNodeAt(unsigned int depth) const
{
Serializer s;
nodeAt(depth)->serializeWithPrefix(s);
return s.modData();
}
/**
* Every node below the root, down to and including the deepest one.
*
* @param firstDepth The shallowest node to include, so a caller can feed
* the chain in more than one batch.
* @return The nodes, each with its claimed position.
*/
[[nodiscard]] std::vector<std::pair<SHAMapNodeID, SHAMapTreeNodePtr>>
nodesBelowRoot(unsigned int firstDepth = 1) const
{
std::vector<std::pair<SHAMapNodeID, SHAMapTreeNodePtr>> data;
for (auto depth = firstDepth; depth <= deepestDepth; ++depth)
data.emplace_back(idAt(depth), nodeAt(depth));
return data;
}
/**
* Fill a synching map, stopping one level short of the deepest node so the
* caller offers that one itself.
*
* Reports rather than asserts: the two binaries sharing this header use
* different test frameworks.
*
* @param map The map to fill.
* @return Whether the root and every node above the deepest one was
* accepted, which is a property of the chain.
*/
[[nodiscard]] bool
fill(SHAMap& map) const
{
if (!map.addRootNode(rootHash, nodeAt(0), nullptr).isGood())
return false;
for (auto depth = 1u; depth < deepestDepth; ++depth)
{
if (!map.addKnownNode(idAt(depth), nodeAt(depth), nullptr).isUseful())
return false;
}
return true;
}
/**
* Offer the deepest node, which for a fabricated chain is the inner node at
* SHAMap::kLeafDepth, a depth only a leaf may occupy.
*
* @param map The map to offer the node to, filled by fill() first.
* @return The verdict addKnownNode() reached.
*/
[[nodiscard]] SHAMapAddNode
addOffendingNode(SHAMap& map) const
{
return map.addKnownNode(idAt(deepestDepth), nodeAt(deepestDepth), nullptr);
}
private:
/**
* Build any of the shapes.
*
* @param leafDepth Where a real transaction leaf sits, or nullopt to run
* inner nodes all the way to SHAMap::kLeafDepth instead.
* @param seed Varies the chain's contents. See the public entry points.
*/
DeepChain(std::optional<unsigned int> leafDepth, unsigned int seed)
: nodes(leafDepth.value_or(SHAMap::kLeafDepth) + 1)
{
if (!leafDepth)
{
// With no leaf depth given, the deepest inner node points at a child that stays
// unresolvable.
buildInnersDownTo(SHAMap::kLeafDepth, SHAMapHash{UInt256{seed}});
return;
}
// Exactly kLeafItemBytes of payload, the smallest a leaf item may be. Checked rather than
// assumed, since a caller relates that constant to a threshold of its own.
Serializer payload;
payload.add32(seed);
payload.add32(0);
payload.add32(0);
if (payload.size() != kLeafItemBytes)
Throw<std::logic_error>("DeepChain: unexpected leaf payload size");
Serializer wire;
wire.addRaw(payload.peekData());
wire.add8(kWireTypeTransaction);
auto const leaf = SHAMapTreeNode::makeFromWire(makeSlice(wire.peekData()));
// A transaction leaf's key is the hash of its own contents, so its position follows
// this key's nibbles.
pathKey = leafKey(*leaf);
deepestDepth = *leafDepth;
nodes[*leafDepth] = leaf;
if (*leafDepth == 0)
{
// The leaf is the root, so the build stops here.
rootHash = leaf->getHash();
return;
}
buildInnersDownTo(*leafDepth - 1, leaf->getHash());
}
/**
* Fill in inner nodes, each with one real child, from the root down to the
* given depth, and record the root hash.
*
* Bottom-up, since each node's hash covers the child hash below it.
*
* @param deepest The depth of the deepest inner node to build. May be
* SHAMap::kLeafDepth, which is the fabricated chain's whole point.
* @param childHash What that deepest inner node points at.
*/
void
buildInnersDownTo(unsigned int deepest, SHAMapHash childHash)
{
for (auto depth = deepest + 1; depth-- > 0;)
{
// A key has 64 nibbles, so SHAMap::kLeafDepth is one past the last nibble
// selectBranch() reads from a 32-byte key. A fabricated chain's pathKey is zero, so
// branch 0 is the position such a node claims.
auto const branch =
depth == SHAMap::kLeafDepth ? 0u : selectBranch(idAt(depth), pathKey);
Serializer s;
s.addBitString(childHash.asUInt256());
s.add8(static_cast<unsigned char>(branch));
s.add8(kWireTypeCompressedInner);
auto node = SHAMapTreeNode::makeFromWire(makeSlice(s.peekData()));
childHash = node->getHash();
nodes[depth] = std::move(node);
}
rootHash = childHash;
}
};
} // namespace xrpl::tests

View File

@@ -0,0 +1,100 @@
#pragma once
#include <xrpl/basics/SHAMapHash.h>
#include <xrpl/basics/Slice.h>
#include <xrpl/basics/base_uint.h>
#include <xrpl/basics/contract.h>
#include <xrpl/protocol/Serializer.h>
#include <xrpl/shamap/SHAMapInnerNode.h>
#include <xrpl/shamap/SHAMapTreeNode.h>
#include <array>
#include <cstdint>
#include <stdexcept>
#include <vector>
namespace xrpl::tests {
/**
* One child of an inner node a test assembles by hand.
*/
struct InnerChild
{
unsigned int branch{};
SHAMapHash hash;
};
/**
* Assemble an inner node in the wire format's full form.
*
* A full inner node is all 16 branch hashes back to back in branch order,
* followed by the wire type byte. A branch absent from `children`, or given a
* zero hash, yields an empty branch, since the parser derives which branches
* exist from which hashes are non-zero. The node's hash is already correct,
* since makeFromWire() computes it.
*
* @param children The branch and hash of each child to record. Throws if a
* branch is at or above SHAMapInnerNode::kBranchFactor, or if one
* branch is named twice.
* @return The node, or nullptr if the bytes do not parse.
*/
[[nodiscard]] inline SHAMapTreeNodePtr
makeFullInnerNode(std::vector<InnerChild> const& children)
{
// The full form carries a slot for every branch, so a child's position comes from the slot it
// is written to.
std::array<UInt256, SHAMapInnerNode::kBranchFactor> hashes{};
// A duplicate branch is refused here, and one past the last to match the compressed parser.
std::uint32_t seen = 0;
for (auto const& child : children)
{
if (child.branch >= SHAMapInnerNode::kBranchFactor)
Throw<std::logic_error>("makeFullInnerNode: branch is past the last one");
auto const bit = 1u << child.branch;
if ((seen & bit) != 0)
Throw<std::logic_error>("makeFullInnerNode: branch named twice");
seen |= bit;
hashes.at(child.branch) = child.hash.asUInt256();
}
Serializer s;
for (auto const& hash : hashes)
s.addBitString(hash);
s.add8(kWireTypeInner);
return SHAMapTreeNode::makeFromWire(makeSlice(s.peekData()));
}
/**
* Assemble an inner node in the wire format's compressed form.
*
* A compressed inner node is one 33-byte chunk per child, each a hash followed
* by the branch it sits on, and then the wire type byte. Its hash is already
* correct, for the same reason makeFullInnerNode()'s is.
*
* Nothing is checked here, unlike in makeFullInnerNode(). The branch travels as
* one byte, so any value up to 255 reaches the parser, which refuses a branch
* at or above SHAMapInnerNode::kBranchFactor and lets a repeated branch
* overwrite the hash recorded for it.
*
* @param children The branch and hash of each child to record.
* @return The node, or nullptr if the bytes do not parse.
*/
[[nodiscard]] inline SHAMapTreeNodePtr
makeCompressedInnerNode(std::vector<InnerChild> const& children)
{
Serializer s;
for (auto const& child : children)
{
s.addBitString(child.hash.asUInt256());
s.add8(static_cast<unsigned char>(child.branch));
}
s.add8(kWireTypeCompressedInner);
return SHAMapTreeNode::makeFromWire(makeSlice(s.peekData()));
}
} // namespace xrpl::tests

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

@@ -1,30 +1,80 @@
#include <xrpl/basics/Blob.h>
#include <xrpl/basics/SHAMapHash.h>
#include <xrpl/basics/Slice.h>
#include <xrpl/basics/base_uint.h>
#include <xrpl/basics/random.h>
#include <xrpl/beast/hash/uhash.h>
#include <xrpl/beast/utility/Journal.h>
#include <xrpl/beast/xor_shift_engine.h>
#include <xrpl/ledger/Ledger.h>
#include <xrpl/protocol/LedgerHeader.h>
#include <xrpl/protocol/Rules.h>
#include <xrpl/protocol/Serializer.h>
#include <xrpl/shamap/SHAMap.h>
#include <xrpl/shamap/SHAMapAddNode.h>
#include <xrpl/shamap/SHAMapItem.h>
#include <xrpl/shamap/SHAMapMissingNode.h>
#include <xrpl/shamap/SHAMapNodeID.h>
#include <xrpl/shamap/SHAMapSyncFilter.h>
#include <xrpl/shamap/SHAMapTreeNode.h>
#include <boost/smart_ptr/intrusive_ptr.hpp>
#include <gtest/gtest.h>
#include <helpers/TestSink.h>
#include <shamap/DeepChain.h>
#include <shamap/InnerNode.h>
#include <shamap/common.h>
#include <chrono>
#include <cstddef>
#include <cstdint>
#include <list>
#include <map>
#include <optional>
#include <unordered_set>
#include <utility>
#include <vector>
namespace xrpl::tests {
// The cap on how many nodes a walk reports, above what any test here expects.
static constexpr int kMaxNodesPerRequest = 2048;
/**
* Rules with no amendments enabled.
*
* @return The rules.
*/
[[nodiscard]] static Rules
noAmendments()
{
return Rules{std::unordered_set<UInt256, beast::Uhash<>>{}};
}
/**
* Whether a verdict carries exactly the given counts.
*
* The counts rather than get(): that string is a log format, not an API. It is
* pinned once, in the SHAMapAddNode tests, and read here only to describe a
* failure.
*
* @param san The verdict to check.
* @param good How many nodes the batch should have hooked in.
* @param bad How many it should have rejected.
* @param duplicate How many it should have already held.
* @return Whether the verdict matches, naming the actual tally if it does not.
*/
[[nodiscard]] static ::testing::AssertionResult
tallyIs(SHAMapAddNode const& san, int good, int bad, int duplicate)
{
if (san.getGood() == good && san.getBad() == bad && san.getDuplicate() == duplicate)
return ::testing::AssertionSuccess();
return ::testing::AssertionFailure() << "tally is " << san.get() << ", expected good:" << good
<< " bad:" << bad << " dupe:" << duplicate;
}
class SHAMapSyncTest : public ::testing::Test
{
protected:
@@ -81,8 +131,344 @@ protected:
return true;
}
/**
* A sync filter that records every node it is told about, and serves back
* only the ones it was explicitly asked to hold.
*
* Serving is opt-in: the sync path consults the filter before deciding
* a node is missing, so each case serves only the node it wants resolved.
*/
class RecordingFilter : public SHAMapSyncFilter
{
public:
// What one gotNode() call was told, in the order the calls arrived.
struct Report
{
bool fromFilter;
SHAMapHash hash;
std::uint32_t ledgerSeq;
};
void
gotNode(
bool fromFilter,
SHAMapHash const& hash,
std::uint32_t ledgerSeq,
Blob&&, // NOLINT(cppcoreguidelines-rvalue-reference-param-not-moved)
SHAMapNodeType) const override
{
reports_.push_back({.fromFilter = fromFilter, .hash = hash, .ledgerSeq = ledgerSeq});
}
[[nodiscard]] std::optional<Blob>
getNode(SHAMapHash const& hash) const override
{
if (auto const it = served_.find(hash); it != served_.end())
return it->second;
return std::nullopt;
}
/**
* Offer a node back to the map, as a fetch pack does.
*
* @param node The node to serve, keyed by its own hash.
*/
void
serve(SHAMapTreeNodePtr const& node)
{
Serializer s;
node->serializeWithPrefix(s);
served_.emplace(node->getHash(), s.modData());
}
[[nodiscard]] std::vector<Report> const&
reports() const
{
return reports_;
}
private:
// Mutable because the whole interface is const: a filter is handed to the map by
// const pointer, so recording has to happen through one.
mutable std::vector<Report> reports_;
std::map<SHAMapHash, Blob> served_;
};
/**
* A root inner node with all 16 branches occupied and not one of them
* resolvable.
*
* A walk of a backed map posts an asynchronous read for every branch in a
* single pass, so the nodestore reader threads run finishFetch() for the
* same map at the same time.
*/
struct WideRoot
{
SHAMapTreeNodePtr node;
SHAMapHash hash;
WideRoot()
{
std::vector<InnerChild> children;
children.reserve(SHAMap::kBranchFactor);
for (auto branch = 0u; branch < SHAMap::kBranchFactor; ++branch)
{
// Derived from the branch, so each posts its own read and no hash collides with a
// real node's.
UInt256 childHash;
childHash.begin()[0] = 0xFA;
childHash.begin()[1] = 0xB1;
childHash.begin()[2] = static_cast<unsigned char>(branch);
children.push_back({.branch = branch, .hash = SHAMapHash{childHash}});
}
node = makeFullInnerNode(children);
hash = node->getHash();
}
};
};
// Only a leaf may sit at kLeafDepth. An inner node there is reported as bad data and leaves the
// map invalid.
TEST_F(SHAMapSyncTest, inner_node_at_leaf_depth)
{
TestNodeFamily f{j_};
DeepChain const chain;
SHAMap map{SHAMapType::FREE, f};
map.setSynching();
ASSERT_TRUE(chain.fill(map));
ASSERT_TRUE(map.isValid());
auto const result = chain.addOffendingNode(map);
EXPECT_TRUE(tallyIs(result, 0, 1, 0));
EXPECT_FALSE(result.isGood());
EXPECT_FALSE(map.isValid());
}
// A node the descent rejects is bad data: the batch counts it bad and the map stays usable for
// another sender. All three ways of getting there are covered, since they share that verdict.
TEST_F(SHAMapSyncTest, node_that_cannot_be_hooked_is_bad_data)
{
TestNodeFamily f{j_};
DeepChain const chain;
SHAMap map{SHAMapType::FREE, f};
map.setSynching();
ASSERT_TRUE(map.addRootNode(chain.rootHash, chain.nodeAt(0), nullptr).isGood());
// nodeAt(1) is the node the root is missing and its hash matches, but we claim depth 2.
auto const wrongDepth = map.addKnownNode(SHAMapNodeID{2, UInt256{}}, chain.nodeAt(1), nullptr);
EXPECT_TRUE(tallyIs(wrongDepth, 0, 1, 0));
EXPECT_FALSE(wrongDepth.isUseful());
// The chain sits on branch 0 at every depth, so a node claiming a position on branch 1 asks the
// descent to follow a branch the root leaves empty.
UInt256 otherBranch;
otherBranch.begin()[0] = 0x10;
auto const emptyBranch =
map.addKnownNode(SHAMapNodeID{1, otherBranch}, chain.nodeAt(1), nullptr);
EXPECT_TRUE(tallyIs(emptyBranch, 0, 1, 0));
// The right position this time, but the data hashes to something other than the child the root
// says belongs there.
auto const corrupt = map.addKnownNode(SHAMapNodeID{1, UInt256{}}, chain.nodeAt(2), nullptr);
EXPECT_TRUE(tallyIs(corrupt, 0, 1, 0));
// The verdict is bad data alone, so the map stays usable.
EXPECT_TRUE(map.isValid());
}
// A root is installed once and a map is synced against one hash, so a root offered under a hash the
// map does not hold names another tree and is bad data. The same root under the hash the map does
// hold is the duplicate it is.
TEST_F(SHAMapSyncTest, add_root_node_judges_the_hash_asked_for)
{
TestNodeFamily f{j_};
DeepChain const chain;
SHAMap map{SHAMapType::FREE, f};
map.setSynching();
ASSERT_TRUE(map.addRootNode(chain.rootHash, chain.nodeAt(0), nullptr).isGood());
// The same root under the hash the map holds: already held, and reported as such.
auto const same = map.addRootNode(chain.rootHash, chain.nodeAt(0), nullptr);
EXPECT_TRUE(tallyIs(same, 0, 0, 1));
EXPECT_TRUE(same.isGood());
// A hash the map does not hold. nodeAt(1) is a real node of the same chain, so this is a
// well-formed hash that names another tree.
auto const other = map.addRootNode(chain.nodeAt(1)->getHash(), chain.nodeAt(0), nullptr);
EXPECT_TRUE(tallyIs(other, 0, 1, 0));
EXPECT_FALSE(other.isGood());
EXPECT_TRUE(other.isInvalid());
// The refusal is about the hash asked for, so the root the map holds stays in place and the map
// stays usable.
EXPECT_EQ(map.getHash(), chain.rootHash);
EXPECT_TRUE(map.isValid());
}
// The verdict outranks the full-below cache. That cache is keyed by node hash and shared by every
// map of a family, and a hash covers a node's children but not its depth, so an earlier walk can
// mark the same subtree hash complete at one depth while this map reaches it at kLeafDepth, with no
// collision involved. The descent therefore skips the lookup at that boundary and reaches the depth
// verdict first. This case seeds the entry a lookup would match, so dropping the skip turns the
// verdict back into a duplicate.
TEST_F(SHAMapSyncTest, map_invalidating_node_is_judged_before_the_cache_is_read)
{
TestNodeFamily f{j_};
DeepChain const chain;
SHAMap map{SHAMapType::FREE, f};
map.setSynching();
ASSERT_TRUE(chain.fill(map));
ASSERT_TRUE(map.isValid());
// This case seeds the entry a lookup at the boundary would match: the offending node's own
// hash. The descent skips the lookup there, so the entry is never read and the depth verdict
// stands. Drop the skip and the hit returns for the whole branch, so the tally below becomes
// a duplicate.
f.getFullBelowCache()->insert(chain.nodeAt(SHAMap::kLeafDepth)->getHash().asUInt256());
auto const result = chain.addOffendingNode(map);
EXPECT_TRUE(tallyIs(result, 0, 1, 0));
EXPECT_FALSE(result.isGood());
EXPECT_FALSE(map.isValid());
}
// A map marked complete in the database withdraws that claim the first time a read misses, and
// reports the miss once so the ledger can be re-acquired. Sixteen unresolvable branches are posted
// in one pass, so with four reader threads the misses overlap and finishFetch() runs concurrently
// for a single map.
//
// This pins the report path: a miss withdraws the flag and reports once. The nightly
// ThreadSanitizer job that PR 8245 adds covers the ordering.
TEST_F(SHAMapSyncTest, full_flag_is_withdrawn_once_by_concurrent_readers)
{
static constexpr auto kRounds = 8uz;
static constexpr auto kReadThreads = 4;
for (auto round = 0uz; round < kRounds; ++round)
{
TestNodeFamily f{j_, kReadThreads};
WideRoot const root;
// Backed, so descendAsync() posts real asynchronous reads rather than resolving inline.
SHAMap map{SHAMapType::FREE, f};
map.setSynching();
ASSERT_TRUE(map.addRootNode(root.hash, root.node, nullptr).isGood());
// The claim the first miss has to withdraw.
map.setFull();
// A null filter, so every branch is read from a database that lacks it.
EXPECT_EQ(map.getMissingNodes(kMaxNodesPerRequest, nullptr).size(), SHAMap::kBranchFactor)
<< "round " << round;
EXPECT_EQ(f.missingBySeqReports(), 1uz) << "round " << round;
}
}
// Every node the sync path hands to a filter carries the map's ledger sequence. All three call
// sites are covered: a root taken from a peer, a node taken from a peer, and a node the walk
// resolved out of the filter itself.
TEST_F(SHAMapSyncTest, sync_filter_is_told_the_ledger_sequence)
{
static constexpr std::uint32_t kLedgerSeq = 7;
TestNodeFamily f{j_};
// A three-level chain, built bottom-up so each hash covers the one below it. Every node keeps
// its one child on branch 0. The deepest points at a child the fixture withholds, so the walk
// always has something to ask for.
auto const deepest = makeCompressedInnerNode({{.branch = 0u, .hash = SHAMapHash{UInt256{1}}}});
auto const middle = makeCompressedInnerNode({{.branch = 0u, .hash = deepest->getHash()}});
auto const root = makeCompressedInnerNode({{.branch = 0u, .hash = middle->getHash()}});
// Unbacked, so the filter is the only source, and a node it withholds counts as missing.
SHAMap map{SHAMapType::FREE, f};
map.setUnbacked();
map.setSynching();
map.setLedgerSeq(kLedgerSeq);
RecordingFilter filter;
ASSERT_TRUE(map.addRootNode(root->getHash(), root, &filter).isGood());
ASSERT_TRUE(map.addKnownNode(SHAMapNodeID{1, UInt256{}}, middle, &filter).isUseful());
// The deepest node becomes resolvable only now, after the two above were added directly.
filter.serve(deepest);
auto const missing = map.getMissingNodes(kMaxNodesPerRequest, &filter);
// The walk resolved the deepest node through the filter and then asked for its child.
ASSERT_EQ(missing.size(), 1u);
EXPECT_EQ(missing[0].second, UInt256{1});
ASSERT_EQ(filter.reports().size(), 3u);
// The two nodes taken from a peer, which the filter is told about so it can store them.
EXPECT_FALSE(filter.reports()[0].fromFilter);
EXPECT_EQ(filter.reports()[0].hash, root->getHash());
EXPECT_FALSE(filter.reports()[1].fromFilter);
EXPECT_EQ(filter.reports()[1].hash, middle->getHash());
// The one the walk read back out of the filter, which is reported as such.
EXPECT_TRUE(filter.reports()[2].fromFilter);
EXPECT_EQ(filter.reports()[2].hash, deepest->getHash());
for (auto const& report : filter.reports())
EXPECT_EQ(report.ledgerSeq, kLedgerSeq) << "hash " << report.hash;
}
// Ledger::setFull() publishes each map's ledger sequence alongside the flag that lets the first
// nodestore miss report a gap. The sequence is what the lookup resolving that gap reads.
//
// This pins that setFull() sets the sequence. The nightly ThreadSanitizer job that PR 8245
// adds covers the order of the two stores.
TEST_F(SHAMapSyncTest, ledger_set_full_publishes_the_ledger_sequence)
{
static constexpr std::uint32_t kLedgerSeq = 7;
TestNodeFamily f{j_};
LedgerHeader header;
header.seq = kLedgerSeq;
// Non-zero, so the map has a root to look for and the lookup can miss.
header.txHash = UInt256{1};
header.hash = calculateLedgerHash(header);
Ledger ledger{header, noAmendments(), f};
// The constructor already looked for that root and missed. The map becomes complete only at
// setFull() below, so the report count is still zero here.
ASSERT_EQ(f.missingBySeqReports(), 0uz);
ledger.setFull();
// Still missing, and now the map has a claim to withdraw, so the gap is reported.
// TestNodeFamily throws in place of the real family's re-acquisition, which finishFetch()
// logs and swallows.
EXPECT_FALSE(ledger.txMap().fetchRoot(SHAMapHash{header.txHash}, nullptr));
EXPECT_EQ(f.missingBySeqReports(), 1uz);
EXPECT_EQ(f.missingBySeqRefNum(), kLedgerSeq);
}
TEST_F(SHAMapSyncTest, sync)
{
TestNodeFamily f{j_}, f2{j_};
@@ -92,7 +478,6 @@ TEST_F(SHAMapSyncTest, sync)
static constexpr auto kItemCount = 10000uz;
static constexpr auto kInvariantInterval = 100uz;
static constexpr auto kNodesToConfuse = 500uz;
static constexpr auto kMaxNodesPerRequest = 2048;
for (auto i = 0uz; i < kItemCount; ++i)
{
@@ -180,4 +565,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

@@ -14,7 +14,9 @@
#include <xrpl/shamap/FullBelowCache.h>
#include <xrpl/shamap/TreeNodeCache.h>
#include <atomic>
#include <chrono>
#include <cstddef>
#include <cstdint>
#include <memory>
#include <stdexcept>
@@ -26,16 +28,27 @@ class TestNodeFamily : public Family
private:
std::unique_ptr<node_store::Database> db_;
// Declared before the two caches, which both bind a reference to it in the initializer list.
TestStopwatch clock_;
std::shared_ptr<FullBelowCache> fbCache_;
std::shared_ptr<TreeNodeCache> tnCache_;
TestStopwatch clock_;
node_store::DummyScheduler scheduler_;
beast::Journal const j_;
// Written from whichever nodestore reader thread reports the miss, so read back atomically.
std::atomic<std::size_t> missingBySeqReports_ = 0;
std::atomic<std::uint32_t> missingBySeqRefNum_ = 0;
public:
TestNodeFamily(beast::Journal j)
/**
* @param j The journal to log through.
* @param readThreads How many nodestore reader threads to run asynchronous
* fetches on.
*/
explicit TestNodeFamily(beast::Journal j, int readThreads = 1)
: fbCache_(std::make_shared<FullBelowCache>("App family full below cache", clock_, j))
, tnCache_(
std::make_shared<TreeNodeCache>(
@@ -50,7 +63,7 @@ public:
testSection.set(Keys::kType, "memory");
testSection.set(Keys::kPath, "SHAMap_test");
db_ = node_store::Manager::instance().makeDatabase(
megabytes(4), scheduler_, 1, testSection, j);
megabytes(4), scheduler_, readThreads, testSection, j);
}
node_store::Database&
@@ -90,14 +103,28 @@ public:
tnCache_->sweep();
}
/**
* Record the report and throw, standing in for Family's real acquisition
* machinery.
*
* @param refNum Sequence of the ledger with the missing node. Recorded, and
* readable through missingBySeqRefNum().
* @param nodeHash Hash of the missing node. Unused.
*/
void
missingNodeAcquireBySeq(
[[maybe_unused]] std::uint32_t refNum,
[[maybe_unused]] UInt256 const& nodeHash) override
missingNodeAcquireBySeq(std::uint32_t refNum, [[maybe_unused]] UInt256 const& nodeHash) override
{
missingBySeqRefNum_.store(refNum, std::memory_order_release);
++missingBySeqReports_;
Throw<std::runtime_error>("missing node");
}
/**
* Throw, standing in for Family's real acquisition machinery. Uncounted.
*
* @param refHash Hash of the ledger with the missing node. Unused.
* @param refNum Sequence of the ledger with the missing node. Unused.
*/
void
missingNodeAcquireByHash(
[[maybe_unused]] UInt256 const& refHash,
@@ -106,6 +133,29 @@ public:
Throw<std::runtime_error>("missing node");
}
/**
* How many times a map of this family has withdrawn its claim of being
* complete in the database. Counted per family, not per map.
*
* @return The number of missingNodeAcquireBySeq() calls so far.
*/
[[nodiscard]] std::size_t
missingBySeqReports() const
{
return missingBySeqReports_.load(std::memory_order_acquire);
}
/**
* The ledger sequence the most recent such report named.
*
* @return The sequence, or zero if nothing has been reported yet.
*/
[[nodiscard]] std::uint32_t
missingBySeqRefNum() const
{
return missingBySeqRefNum_.load(std::memory_order_acquire);
}
void
reset() override
{

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;
@@ -155,6 +186,17 @@ private:
bool
takeHeader(std::string_view data);
/**
* Fail the acquisition when the header's account hash is zero. No ledger
* has an empty state map, so such a header cannot name a ledger. Both
* tryDB() and takeHeader() judge the header here.
*
* @return Whether the acquisition was failed. Then failed_ is set and
* ledger_ is null.
*/
bool
failOnZeroAccountHash();
void
receiveNode(
std::shared_ptr<Peer> const& peer,

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)
@@ -97,8 +93,12 @@ InboundLedger::init(ScopedLockType& collectionLock)
collectionLock.unlock();
tryDB(app_.getNodeFamily().db());
// done() is what wakes whatever is waiting and records the hash in recentFailures_.
if (failed_)
{
done();
return;
}
if (!complete_)
{
@@ -241,7 +241,12 @@ InboundLedger::tryDB(node_store::Database& srcDB)
<< "hash " << hash_ << " seq " << std::to_string(seq_) << " cannot be a ledger";
ledger_.reset();
failed_ = true;
return;
}
// A zero account hash cannot be a ledger either. The helper sets failed_ and drops the
// ledger, so the callers below store and publish nothing for such a header.
failOnZeroAccountHash();
};
// Try to fetch the ledger header from the DB
@@ -309,12 +314,6 @@ InboundLedger::tryDB(node_store::Database& srcDB)
if (!haveState_)
{
if (ledger_->header().accountHash.isZero())
{
JLOG(journal_.fatal()) << "We are acquiring a ledger with a zero account hash";
failed_ = true;
return;
}
AccountStateSF filter(ledger_->stateMap().family().db(), app_.getLedgerMaster());
if (ledger_->stateMap().fetchRoot(SHAMapHash{ledger_->header().accountHash}, &filter))
{
@@ -497,6 +496,7 @@ InboundLedger::trigger(std::shared_ptr<Peer> const& peer, TriggerReason reason)
if (failed_)
{
JLOG(journal_.warn()) << " failed local for " << hash_;
done();
return;
}
}
@@ -774,15 +774,31 @@ InboundLedger::filterNodes(
recentNodes_.insert(n.second);
}
bool
InboundLedger::failOnZeroAccountHash()
{
bool const zero = ledger_->header().accountHash.isZero();
SOMETIMES(zero, "xrpl::InboundLedger::failOnZeroAccountHash : zero account hash");
if (!zero)
return false;
JLOG(journal_.fatal()) << "We are acquiring a ledger with a zero account hash: " << hash_;
failed_ = true;
ledger_.reset();
return true;
}
/**
* Take ledger header data
* Call with a lock
* Build the ledger from a header a peer supplied. Call with mtx_ held.
*
* @param data The serialized header, without a hash prefix.
* @return False for a header other than the one asked for. True otherwise,
* also when the header cannot name a ledger, and then failed_ is set,
* ledger_ is null, haveHeader_ stays false, and the caller calls done().
*/
// data must not have hash prefix
bool
InboundLedger::takeHeader(std::string_view data)
{
// Return value: true=normal, false=bad data
JLOG(journal_.trace()) << "got header acquiring ledger " << hash_;
if (complete_ || failed_ || haveHeader_)
@@ -798,6 +814,12 @@ InboundLedger::takeHeader(std::string_view data)
ledger_.reset();
return false;
}
// The header is the one asked for, so the peer is not charged and the caller reads failed_.
// The helper drops the ledger, and haveHeader_ stays false, as tryDB() leaves them.
if (failOnZeroAccountHash())
return true;
if (seq_ == 0)
seq_ = ledger_->header().seq;
ledger_->stateMap().setLedgerSeq(seq_);
@@ -812,9 +834,6 @@ InboundLedger::takeHeader(std::string_view data)
if (ledger_->header().txHash.isZero())
haveTransactions_ = true;
if (ledger_->header().accountHash.isZero())
haveState_ = true;
ledger_->txMap().setSynching();
ledger_->stateMap().setSynching();
@@ -1099,6 +1118,13 @@ InboundLedger::processData(std::shared_ptr<Peer> peer, protocol::TMLedgerData co
return -1;
}
// takeHeader() failed the acquisition. Nothing else signals that after isDone().
if (failed_)
{
done();
return 0;
}
san.incUseful();
}

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();