mirror of
https://github.com/XRPLF/rippled.git
synced 2026-10-10 21:58:03 +00:00
Compare commits
6 Commits
mvadari/re
...
bthomee/sh
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
dcd44fa0ae | ||
|
|
a0fd44cc7b | ||
|
|
516ed4d433 | ||
|
|
a790fb3527 | ||
|
|
0514203d03 | ||
|
|
db21f68086 |
@@ -103,6 +103,7 @@ words:
|
||||
- dearmor
|
||||
- decryptor
|
||||
- dedented
|
||||
- dedup
|
||||
- deleteme
|
||||
- demultiplexer
|
||||
- deserializaton
|
||||
@@ -344,6 +345,7 @@ words:
|
||||
- unambiguity
|
||||
- unauthorizes
|
||||
- unauthorizing
|
||||
- undeserializable
|
||||
- unergonomic
|
||||
- unfetched
|
||||
- unfindable
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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:
|
||||
/**
|
||||
@@ -375,16 +401,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 +563,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 +808,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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -555,7 +555,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();
|
||||
@@ -623,7 +623,7 @@ SHAMap::addKnownNode(
|
||||
(treeNode->isInner() && currNodeID.getDepth() == kLeafDepth))
|
||||
{
|
||||
// Map is provably invalid
|
||||
state_ = SHAMapState::Invalid;
|
||||
setInvalid();
|
||||
return SHAMapAddNode::useful();
|
||||
}
|
||||
|
||||
@@ -647,7 +647,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();
|
||||
|
||||
365
src/test/app/AcquireTestHelpers.h
Normal file
365
src/test/app/AcquireTestHelpers.h
Normal 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
|
||||
537
src/test/app/InboundLedger_test.cpp
Normal file
537
src/test/app/InboundLedger_test.cpp
Normal 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
|
||||
490
src/test/app/TransactionAcquire_test.cpp
Normal file
490
src/test/app/TransactionAcquire_test.cpp
Normal 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
|
||||
@@ -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
|
||||
@@ -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()};
|
||||
|
||||
268
src/tests/libxrpl/shamap/DeepChain.h
Normal file
268
src/tests/libxrpl/shamap/DeepChain.h
Normal 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
|
||||
100
src/tests/libxrpl/shamap/InnerNode.h
Normal file
100
src/tests/libxrpl/shamap/InnerNode.h
Normal 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
|
||||
76
src/tests/libxrpl/shamap/SHAMapAddNode.cpp
Normal file
76
src/tests/libxrpl/shamap/SHAMapAddNode.cpp
Normal 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
|
||||
@@ -1,30 +1,55 @@
|
||||
#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/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/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<>>{}};
|
||||
}
|
||||
|
||||
class SHAMapSyncTest : public ::testing::Test
|
||||
{
|
||||
protected:
|
||||
@@ -81,8 +106,224 @@ 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();
|
||||
}
|
||||
};
|
||||
};
|
||||
|
||||
// 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 +333,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 +420,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
|
||||
|
||||
@@ -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
|
||||
{
|
||||
|
||||
@@ -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";
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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)
|
||||
{
|
||||
|
||||
@@ -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();
|
||||
|
||||
|
||||
Reference in New Issue
Block a user