Compare commits

...

8 Commits

Author SHA1 Message Date
Bart
88ad25eb0e fix: Refuse to make an invalid SHAMap or Ledger immutable
SHAMap::setImmutable() returns [[nodiscard]] bool and refuses a map
already proven impossible. trySetState() and clearSynching() are the
only writers of state_ past construction, and Invalid outranks and
survives every other state. clearSynching() moves only a Synching map,
so a walk that ends after the ledger settled leaves the map Immutable.

Ledger::setImmutable() and setAccepted() report the same. They check
mapsValid() first and settle each map on its own, so neither is left
mid-sync because the other refused. The hashes are read before the maps
are settled, since getHash() can unshare a dirty tree, and written to
the header only once both have made it. A locally built or loaded ledger
treats false as a broken invariant. One assembled from peer data
recovers, withdrawing complete_ beside the failure. done() settles the
ledger ahead of its reason switch, so a HISTORY ledger reaches
onLedgerFetched() only once settled.

trigger() now calls done() under mtx_, as every other caller already
did, so the flags done() writes are published under one lock.

InboundLedger::getLedger() reports nothing once the acquisition has
failed. A failure found once the header is held, as trigger() reports
for an invalid map, leaves ledger_ populated, and the loadOldLedger
fallbacks read getLedger() without testing isFailed(), so the guard sits
at the one accessor rather than at each caller. ledger_ itself is kept,
since getJson() still reports on the partial maps.
2026-10-04 06:35:38 -04:00
Bart
ccd8838d09 fix: Report a node the sync path refuses as invalid data
addKnownNode() returns invalid() for the two shapes it declines to hook
in: an inner node at kLeafDepth, a depth only a leaf may occupy, and a
node whose claimed ID is not where the hash-verified descent stopped.
The depth case also marks the map Invalid, which is terminal, since
every node above it hash-verified and so the requested root hash itself
commits to a shape no valid tree has. An ID mismatch leaves the map
sound. The depth test precedes the full-below cache lookup, which is
keyed by hash and so answers for no particular depth.

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

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

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

AcquireTestHelpers.h is the first src/test file to include one from
src/tests, so ordering.txt records a test.app > tests.libxrpl edge. That
direction is deliberate: src/test goes away as its suites migrate, so an
edge into the surviving tree outlives the migration. Nothing under
src/tests includes src/test.
2026-10-04 06:35:36 -04:00
Bart
0514203d03 refactor: Open the acquisition seams its tests need
TransactionAcquire and InboundLedger are no longer final and take a
retry interval defaulted to the constant each used before, so a test
subclass can drive a whole timeout chain in a fraction of a second.
TimeoutCounter's constructor already asserts an interval above 10ms and
below 30s, and every production caller takes the default.
TransactionAcquire::map_ and InboundLedger's trigger(), done() and
TriggerReason are protected, which is the only way a subclass reaches
them.

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

View File

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

View File

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

View File

@@ -257,13 +257,42 @@ public:
header_.validated = true;
}
void
/**
* Mark this ledger as accepted and attempt to make it immutable.
*
* The close-time fields are recorded first, since the ledger hash covers
* them.
*
* @param closeTime The consensus-agreed close time.
* @param closeResolution The close time resolution.
* @param correctCloseTime Whether consensus agreed on the close time. If
* false, kSLcfNoConsensusTime is recorded in closeFlags instead.
* @return What setImmutable() returned.
*/
[[nodiscard]] bool
setAccepted(
NetClock::time_point closeTime,
NetClock::duration closeResolution,
bool correctCloseTime);
void
/**
* Mark this ledger as immutable, so it can no longer be modified.
*
* A ledger built or loaded locally cannot have an invalid map, since only
* a map syncing against hashes from outside can be proven impossible (see
* SHAMap::addKnownNode), so a caller on such a path may treat a false
* return as a broken internal invariant. A caller assembling a ledger from
* peer data may not: for it, false is an outcome a peer can produce.
*
* @param rehash Whether to recompute the ledger hash from the header
* fields. The transaction and account hashes are recomputed from the
* maps too, but only the first time.
* @return false if either map is Invalid, leaving the immutable flag as it
* was and the header untouched. A map invalidated partway through can
* still leave the other one immutable, so a false return means the
* ledger must be discarded rather than retried.
*/
[[nodiscard]] bool
setImmutable(bool rehash = true);
bool
@@ -272,23 +301,33 @@ public:
return immutable_;
}
/* Mark this ledger as "should be full".
/**
* @return Whether both maps can still be the maps the header names. See
* SHAMap::isValid().
*/
[[nodiscard]] bool
mapsValid() const
{
return txMap_.isValid() && stateMap_.isValid();
}
"Full" is metadata property of the ledger, it indicates
that the local server wants all the corresponding nodes
in durable storage.
This is marked `const` because it reflects metadata
and not data that is in common with other nodes on the
network.
*/
/**
* Mark this ledger as "should be full", indicating that the local server
* wants all the corresponding nodes in durable storage.
*
* Const because it reflects metadata, not data this ledger shares with
* other nodes on the network.
*/
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
@@ -418,6 +457,22 @@ private:
static std::pair<std::shared_ptr<STTx const>, std::shared_ptr<STObject const>>
deserializeTxPlusMeta(SHAMapItem const& item);
/**
* Make both maps immutable.
*
* Each map is settled on its own, so neither is left mid-sync because the
* other refused.
*
* @return Whether both maps are immutable.
*/
[[nodiscard]] bool
setMapsImmutable()
{
bool const txImmutable = txMap_.setImmutable();
bool const stateImmutable = stateMap_.setImmutable();
return txImmutable && stateImmutable;
}
bool immutable_;
// A SHAMap containing the transactions associated with this ledger.

View File

@@ -2,6 +2,7 @@
#include <xrpl/basics/Blob.h>
#include <xrpl/basics/IntrusivePointer.h>
#include <xrpl/basics/Log.h>
#include <xrpl/basics/SHAMapHash.h>
#include <xrpl/basics/base_uint.h>
#include <xrpl/beast/utility/Journal.h>
@@ -17,6 +18,7 @@
#include <xrpl/shamap/SHAMapNodeID.h>
#include <xrpl/shamap/SHAMapTreeNode.h>
#include <atomic>
#include <condition_variable>
#include <cstddef>
#include <cstdint>
@@ -40,7 +42,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 +122,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:
/**
@@ -152,7 +179,14 @@ public:
SHAMap&
operator=(SHAMap const&) = delete;
// Take a snapshot of the given map:
/**
* Take a snapshot of the given map.
*
* @param other The map to snapshot. An Invalid source yields an Invalid
* snapshot.
* @param isMutable Whether the snapshot may be modified. Ignored when other
* is Invalid.
*/
SHAMap(SHAMap const& other, bool isMutable);
// build new map
@@ -190,8 +224,15 @@ public:
//--------------------------------------------------------------------------
// Returns a new map that's a snapshot of this one.
// Handles copy on write for mutable snapshots.
/**
* Return a new map that is a snapshot of this one.
*
* Handles copy on write for mutable snapshots. An invalid map yields an
* invalid snapshot.
*
* @param isMutable Whether the snapshot may be modified.
* @return The snapshot.
*/
std::shared_ptr<SHAMap>
snapShot(bool isMutable) const;
@@ -341,6 +382,9 @@ public:
* This function is used when receiving the root node of a SHAMap from a peer during ledger
* synchronization. The node must already have been deserialized.
*
* A root offered under a hash the map does not hold names another tree and
* is reported as invalid data. A root matching that hash is a duplicate.
*
* @param hash The expected hash of the root node.
* @param rootNode A deserialized root node to add.
* @param filter Optional sync filter to track received nodes.
@@ -360,6 +404,10 @@ public:
* is inserted at the position specified by nodeID. The node must already have been
* deserialized.
*
* A node that no valid tree can hold makes the map Invalid, which is
* terminal: the root hash committed to an impossible shape, so no peer
* can satisfy it.
*
* @param nodeID The position in the tree where this node belongs.
* @param treeNode A deserialized tree node to add.
* @param filter Optional sync filter to track received nodes.
@@ -375,16 +423,57 @@ public:
SHAMapTreeNodePtr treeNode,
SHAMapSyncFilter const* filter);
// status functions
void
/**
* Mark this map as immutable, so it can no longer be modified.
*
* @return false if the map is Invalid and was left unchanged, true
* otherwise.
*/
[[nodiscard]] bool
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;
/**
* @return Whether the map is settled, with its hash and nodes fixed.
*/
[[nodiscard]] bool
isImmutable() const;
/**
* Mark this map as syncing, fixing its hash while still allowing missing
* nodes to be added.
*
* The map must be freshly constructed. The body asserts that the map is not
* Invalid, so a caller that ignores the precondition stops a build with
* assertions enabled.
*/
void
setSynching();
/**
* Mark this map as no longer syncing, so it can be modified again.
*
* Moves only a Synching map. An Immutable map stays settled, a Modifying
* map has nothing to clear, and Invalid is terminal.
*/
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 +613,51 @@ 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. Cannot fail, since
* Invalid outranks every other state; see trySetState().
*/
void
setInvalid();
/**
* Move the map to a new state, atomically.
*
* With clearSynching(), the only writer of state_ past construction, so
* the order between the states lives in one place: Invalid outranks all
* of them and is always stored, while every other transition is refused
* once the map is Invalid, which is what makes that verdict terminal.
* clearSynching() is narrower and moves only Synching.
*
* @param desired The state to move to.
* @return false if the map is Invalid and the requested state is not,
* leaving it unchanged; true otherwise.
*/
bool
trySetState(SHAMapState desired);
// tree node cache operations
SHAMapTreeNodePtr
cacheLookup(SHAMapHash const& hash) const;
@@ -741,44 +875,117 @@ 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 void
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 bool
SHAMap::trySetState(SHAMapState desired)
{
// Invalid is stored outright: a walk reaching that verdict has to win against a thread
// settling the map.
if (desired == SHAMapState::Invalid)
{
state_.store(SHAMapState::Invalid, std::memory_order_release);
return true;
}
// A failed exchange both reports the state and refreshes expected, so no load is needed
// ahead of the loop.
auto expected = SHAMapState::Modifying;
while (expected != SHAMapState::Invalid)
{
if (state_.compare_exchange_weak(
expected, desired, std::memory_order_acq_rel, std::memory_order_acquire))
{
return true;
}
}
return false;
}
inline bool
SHAMap::setImmutable()
{
XRPL_ASSERT(state_ != SHAMapState::Invalid, "xrpl::SHAMap::setImmutable : state is valid");
state_ = SHAMapState::Immutable;
SOMETIMES(!isValid(), "xrpl::SHAMap::setImmutable : map is invalid");
return trySetState(SHAMapState::Immutable);
}
inline bool
SHAMap::isSynching() const
{
return state_ == SHAMapState::Synching;
return state() == SHAMapState::Synching;
}
inline bool
SHAMap::isImmutable() const
{
return state() == SHAMapState::Immutable;
}
inline void
SHAMap::setSynching()
{
state_ = SHAMapState::Synching;
// Guarded, so this is not a way out of Invalid.
if (!trySetState(SHAMapState::Synching))
{
// Only ever called on a freshly constructed map.
// LCOV_EXCL_START
UNREACHABLE("xrpl::SHAMap::setSynching : map is invalid");
// LCOV_EXCL_STOP
}
}
inline void
SHAMap::clearSynching()
{
state_ = SHAMapState::Modifying;
// Only Synching moves. A walk can end after the ledger settled, and then the map stays
// Immutable; a Modifying map has nothing to clear; an invalid map stays invalid. Only that
// last refusal is logged, since a concurrent walk is what produces it.
auto expected = SHAMapState::Synching;
if (state_.compare_exchange_strong(
expected, SHAMapState::Modifying, std::memory_order_acq_rel, std::memory_order_acquire))
{
return;
}
bool const invalid = expected == SHAMapState::Invalid;
SOMETIMES(invalid, "xrpl::SHAMap::clearSynching : map is invalid");
if (invalid)
{
JLOG(journal_.warn()) << "Refused to clear synching on an invalid map, root hash "
<< root_->getHash();
}
}
inline bool
SHAMap::isValid() const
{
return state_ != SHAMapState::Invalid;
return state() != SHAMapState::Invalid;
}
inline void
SHAMap::setInvalid()
{
// Through trySetState() like every other transition, so state_ has one writer funnel.
trySetState(SHAMapState::Invalid);
}
inline void

View File

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

View File

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

View File

@@ -202,7 +202,13 @@ Ledger::Ledger(
}
stateMap_.flushDirty(NodeObjectType::AccountNode);
setImmutable();
// Built locally. See Ledger::setImmutable().
if (!setImmutable())
{
// LCOV_EXCL_START
logicError("Ledger::Ledger(CreateGenesisT, ...): genesis ledger map is invalid");
// LCOV_EXCL_STOP
}
}
Ledger::Ledger(
@@ -213,7 +219,7 @@ Ledger::Ledger(
Fees const& fees,
Family& family,
beast::Journal j)
: immutable_(true)
: immutable_(false)
, txMap_(SHAMapType::TRANSACTION, info.txHash, family)
, stateMap_(SHAMapType::STATE, info.accountHash, family)
, fees_(fees)
@@ -236,8 +242,20 @@ Ledger::Ledger(
JLOG(j.warn()) << "Don't have state data root for ledger" << header_.seq;
}
txMap_.setImmutable();
stateMap_.setImmutable();
// Loaded locally. See Ledger::setImmutable().
if (setMapsImmutable())
{
immutable_ = true;
}
else
{
// LCOV_EXCL_START
JLOG(j.error()) << "Invalid map for ledger " << header_.seq;
UNREACHABLE("xrpl::Ledger::Ledger(LedgerHeader const&, ...) : map is invalid");
// Treat it as a damaged ledger: the code below recomputes the hash and re-acquires.
loaded = false;
// LCOV_EXCL_STOP
}
if (!setup())
loaded = false;
@@ -308,27 +326,50 @@ Ledger::Ledger(
setup();
}
void
bool
Ledger::setImmutable(bool rehash)
{
// Force update, since this is the only
// place the hash transitions to valid
if (!immutable_ && rehash)
// A map found structurally invalid during sync must never be made immutable: isValid() tests
// only for Invalid, and an immutable ledger is treated as persistable. Asked before anything is
// written, so a refusal leaves the header exactly as it was rather than half relabeled.
if (!mapsValid())
return false;
// Read here but written to the header below, once the maps are settled: getHash() can unshare a
// dirty tree, so it has to run while the map is still mutable, while a write to the header must
// wait until both maps have made it. Skipped once the ledger is immutable, since its maps can
// no longer change.
bool const deriveMapHashes = !immutable_ && rehash;
UInt256 const txHash = deriveMapHashes ? txMap_.getHash().asUInt256() : UInt256{};
UInt256 const accountHash = deriveMapHashes ? stateMap_.getHash().asUInt256() : UInt256{};
// Both were valid at the check above, but a concurrent walk can invalidate one in between (see
// SHAMap::state_), so the result is checked. setInvalid() outranks Immutable and can land
// after both maps have been settled, so this narrows the window it guards.
bool const bothImmutable = setMapsImmutable();
SOMETIMES(!bothImmutable, "xrpl::Ledger::setImmutable : map invalidated while going immutable");
if (!bothImmutable)
return false;
// Written only now, so losing the race above leaves the header describing what the ledger was
// built from rather than a map that has since been abandoned. Forced rather than conditional,
// since this is the only place the hash transitions to valid.
if (deriveMapHashes)
{
header_.txHash = txMap_.getHash().asUInt256();
header_.accountHash = stateMap_.getHash().asUInt256();
header_.txHash = txHash;
header_.accountHash = accountHash;
}
if (rehash)
header_.hash = calculateLedgerHash(header_);
// Set last, so isImmutable() reports only a ledger whose maps are both immutable.
immutable_ = true;
txMap_.setImmutable();
stateMap_.setImmutable();
setup();
return true;
}
void
bool
Ledger::setAccepted(
NetClock::time_point closeTime,
NetClock::duration closeResolution,
@@ -340,7 +381,17 @@ Ledger::setAccepted(
header_.closeTime = closeTime;
header_.closeTimeResolution = closeResolution;
header_.closeFlags = correctCloseTime ? 0 : kSLcfNoConsensusTime;
setImmutable();
// Built locally. See Ledger::setImmutable().
if (!setImmutable())
{
// LCOV_EXCL_START
JLOG(j_.error()) << "Invalid map for accepted ledger " << header_.seq;
UNREACHABLE("xrpl::Ledger::setAccepted : map is invalid");
return false;
// LCOV_EXCL_STOP
}
return true;
}
bool

View File

@@ -26,6 +26,7 @@
#include <boost/smart_ptr/intrusive_ptr.hpp>
#include <atomic>
#include <cstdint>
#include <exception>
#include <functional>
@@ -77,14 +78,23 @@ 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_)
{
// A snapshot shares the source's root, so Invalid carries over. Read once into a local, since
// the source's state can change while this runs.
auto const otherState = other.state();
auto const ownState = [&] {
if (otherState == SHAMapState::Invalid)
return SHAMapState::Invalid;
return isMutable ? SHAMapState::Modifying : SHAMapState::Immutable;
}();
state_.store(ownState, std::memory_order_release);
// If either map may change, they cannot share nodes
if ((state_ != SHAMapState::Immutable) || (other.state_ != SHAMapState::Immutable))
if ((ownState != SHAMapState::Immutable) || (otherState != SHAMapState::Immutable))
{
unshare();
}
@@ -105,7 +115,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 +180,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 +193,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 +233,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 +413,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 +438,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 +651,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 +733,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 +824,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 +915,7 @@ SHAMap::writeNode(NodeObjectType t, SHAMapTreeNodePtr node) const
Serializer s;
node->serializeWithPrefix(s);
f_.db().store(t, std::move(s.modData()), node->getHash().asUInt256(), ledgerSeq_);
f_.db().store(t, std::move(s.modData()), node->getHash().asUInt256(), ledgerSeq());
return node;
}

View File

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

View File

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

View File

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

View File

@@ -0,0 +1,666 @@
#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/ledger/Ledger.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);
}
/**
* Record that every part has been fetched.
*/
void
markComplete()
{
ScopedLockType const sl(mtx_);
complete_ = true;
}
/**
* Settle the acquisition and signal whatever is waiting on it.
*/
void
signalDone()
{
ScopedLockType const sl(mtx_);
done();
}
};
/**
* The ledger an acquisition is assembling, as a pointer that can modify it.
*
* @param acquire The acquisition to read from.
* @return The ledger, or nullptr if there is none to report.
*/
[[nodiscard]] static std::shared_ptr<Ledger>
mutableLedger(InboundLedger const& acquire)
{
// Sound because the acquisition holds a non-const ledger and only hands out a const view.
return std::const_pointer_cast<Ledger>(acquire.getLedger());
}
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.
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 reason decides where a settled ledger goes: HISTORY counts it toward the fetch rate
// instead of handing it to LedgerMaster.
{
auto const historyChain = DeepChain::toLeaf(2, nextSeed());
auto const historyHeader = makeHeader(historyChain);
storeHeader(env, historyHeader);
storeStateNodes(env, historyHeader, historyChain, historyChain.deepestDepth);
// onLedgerFetched() is the only writer of this rate, and done() calls it before it
// posts any work, so the read below is not racing the job queue.
auto const rateBefore = env.app().getInboundLedgers().fetchRate();
auto history = std::make_shared<InboundLedger>(
env.app(),
historyHeader.hash,
historyHeader.seq,
InboundLedger::Reason::HISTORY,
stopwatch(),
std::make_unique<RequestCountingPeerSet>());
BEAST_EXPECT(history->checkLocal());
BEAST_EXPECT(history->isComplete());
BEAST_EXPECT(!history->isFailed());
// The switch the rate check reads is reached only once the ledger has been settled.
auto const historyLedger = history->getLedger();
BEAST_EXPECT(historyLedger != nullptr);
if (historyLedger)
BEAST_EXPECT(historyLedger->isImmutable());
BEAST_EXPECT(env.app().getInboundLedgers().fetchRate() > rateBefore);
}
// The same ledger through checkLocal(), which unlike init() reaches done().
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));
}
/**
* A ledger whose map goes invalid on the way to being settled must be
* discarded rather than delivered.
*
* done() settles the ledger before it logs or acts on the outcome, and an
* abandoned map makes settling refuse, so the acquisition records a
* failure.
*
* @param env The environment to run in.
*/
void
testInvalidatedLedgerFailsInDone(jtx::Env& env)
{
testcase("A ledger invalidated on its way to being settled fails");
// The fabricated chain, so feeding it to the state map invalidates the map.
DeepChain const chain{nextSeed()};
// Only the header is local, so the acquisition holds a ledger with an empty state map.
auto const header = makeHeader(chain);
storeHeader(env, header);
auto acquire = std::make_shared<TestableInboundLedger>(
env.app(),
header.hash,
header.seq,
InboundLedger::Reason::GENERIC,
stopwatch(),
std::make_unique<RequestCountingPeerSet>());
BEAST_EXPECT(!acquire->checkLocal());
BEAST_EXPECT(!acquire->isFailed());
BEAST_EXPECT(!acquire->isComplete());
auto const ledger = mutableLedger(*acquire);
BEAST_EXPECT(ledger != nullptr);
if (!ledger)
return;
// The state of affairs done() is handed: every part fetched, as far as the caller can tell.
acquire->markComplete();
// And the walk that has since reached the verdict.
auto& stateMap = ledger->stateMap();
BEAST_EXPECT(stateMap.addRootNode(chain.rootHash, chain.nodeAt(0), nullptr).isGood());
for (auto const& [nodeID, node] : chain.nodesBelowRoot())
stateMap.addKnownNode(nodeID, node, nullptr);
BEAST_EXPECT(!stateMap.isValid());
acquire->signalDone();
// complete_ is withdrawn alongside the failure, or every guard that checks it before
// failed_ keeps treating this ledger as delivered.
BEAST_EXPECT(!acquire->isComplete());
BEAST_EXPECT(acquire->isFailed());
// getLedgerByHash answers null for the hash, and the hash is remembered as a failure,
// which defers re-acquisition.
BEAST_EXPECT(env.app().getLedgerMaster().getLedgerByHash(header.hash) == nullptr);
BEAST_EXPECT(waitFor([&] { return 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);
testInvalidatedLedgerFailsInDone(env);
testLocalFailureSignalsDone(env);
testPeerZeroAccountHashFails(env);
testPeerHeaderWithoutTransactionsCompletes(env);
// Last: the only case that waits out a whole timeout chain.
testTimerRetriesThenGivesUp(env);
}
private:
unsigned int seed_{0};
};
BEAST_DEFINE_TESTSUITE(InboundLedger, app, xrpl);
} // namespace xrpl::test

View File

@@ -11,6 +11,7 @@
#include <xrpl/basics/base_uint.h>
#include <xrpl/basics/chrono.h>
#include <xrpl/basics/contract.h>
#include <xrpl/beast/insight/NullCollector.h>
#include <xrpl/beast/unit_test/suite.h>
#include <xrpl/ledger/ApplyView.h>
@@ -22,6 +23,7 @@
#include <cassert>
#include <memory>
#include <stdexcept>
#include <vector>
namespace xrpl::test {
@@ -70,11 +72,15 @@ public:
}
res->unshare();
// Accept ledger
res->setAccepted(
res->header().closeTime,
res->header().closeTimeResolution,
true /* close time correct*/);
// Accept the ledger. Thrown rather than asserted because a refusal leaves res unusable,
// and this helper is static, so BEAST_EXPECT is out of reach.
if (!res->setAccepted(
res->header().closeTime,
res->header().closeTimeResolution,
true /* close time correct*/))
{
Throw<std::runtime_error>("makeLedger: ledger could not be accepted");
}
lh.insert(res, false);
return res;
}

View File

@@ -100,7 +100,7 @@ class RCLValidations_test : public beast::unit_test::Suite
BEAST_EXPECT(next->read(keylet::feeSettings()));
if (forceHash)
{
next->setImmutable();
BEAST_EXPECT(next->setImmutable());
forceHash = false;
}

View File

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

View File

@@ -13,6 +13,7 @@
#include <xrpl/basics/base_uint.h>
#include <xrpl/basics/chrono.h>
#include <xrpl/basics/contract.h>
#include <xrpl/beast/hash/uhash.h>
#include <xrpl/beast/unit_test/suite.h>
#include <xrpl/json/json_value.h>
@@ -190,7 +191,7 @@ public:
metadata->add(*metaSerializer);
ledger->rawTxInsert(UInt256{1}, txSerializer, metaSerializer);
ledger->setImmutable();
BEAST_EXPECT(ledger->setImmutable());
ledger->setValidated();
try
@@ -260,7 +261,12 @@ public:
ledger->rawTxInsert(UInt256{seq}, txSerializer, metaSerializer);
}
ledger->setImmutable();
// Thrown rather than asserted because a refusal leaves the ledger unusable, and this
// helper is static, so BEAST_EXPECT is out of reach.
if (!ledger->setImmutable())
{
Throw<std::runtime_error>("bookChangesFor: ledger could not be made immutable");
}
ledger->setValidated();
return xrpl::rpc::computeBookChanges(std::static_pointer_cast<Ledger const>(ledger));

View File

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

View File

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

View File

@@ -207,7 +207,10 @@ TxTest::close()
accum.apply(*newLedger);
}
newLedger->setAccepted(ledgerCloseTime, newLedger->header().closeTimeResolution, true);
if (!newLedger->setAccepted(ledgerCloseTime, newLedger->header().closeTimeResolution, true))
{
Throw<std::runtime_error>("TxTest::close: ledger has an invalid map");
}
closedLedger_ = newLedger;

View File

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

View File

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

View File

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

View File

@@ -1,30 +1,82 @@
#include <xrpl/basics/Blob.h>
#include <xrpl/basics/SHAMapHash.h>
#include <xrpl/basics/Slice.h>
#include <xrpl/basics/base_uint.h>
#include <xrpl/basics/chrono.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/Fees.h>
#include <xrpl/protocol/LedgerHeader.h>
#include <xrpl/protocol/Rules.h>
#include <xrpl/protocol/Serializer.h>
#include <xrpl/shamap/SHAMap.h>
#include <xrpl/shamap/SHAMapAddNode.h>
#include <xrpl/shamap/SHAMapItem.h>
#include <xrpl/shamap/SHAMapMissingNode.h>
#include <xrpl/shamap/SHAMapNodeID.h>
#include <xrpl/shamap/SHAMapSyncFilter.h>
#include <xrpl/shamap/SHAMapTreeNode.h>
#include <boost/smart_ptr/intrusive_ptr.hpp>
#include <gtest/gtest.h>
#include <helpers/TestSink.h>
#include <shamap/DeepChain.h>
#include <shamap/InnerNode.h>
#include <shamap/common.h>
#include <chrono>
#include <cstddef>
#include <cstdint>
#include <list>
#include <map>
#include <optional>
#include <unordered_set>
#include <utility>
#include <vector>
namespace xrpl::tests {
// The cap on how many nodes a walk reports, above what any test here expects.
static constexpr int kMaxNodesPerRequest = 2048;
/**
* Rules with no amendments enabled.
*
* @return The rules.
*/
[[nodiscard]] static Rules
noAmendments()
{
return Rules{std::unordered_set<UInt256, beast::Uhash<>>{}};
}
/**
* Whether a verdict carries exactly the given counts.
*
* The counts rather than get(): that string is a log format, not an API. It is
* pinned once, in the SHAMapAddNode tests, and read here only to describe a
* failure.
*
* @param san The verdict to check.
* @param good How many nodes the batch should have hooked in.
* @param bad How many it should have rejected.
* @param duplicate How many it should have already held.
* @return Whether the verdict matches, naming the actual tally if it does not.
*/
[[nodiscard]] static ::testing::AssertionResult
tallyIs(SHAMapAddNode const& san, int good, int bad, int duplicate)
{
if (san.getGood() == good && san.getBad() == bad && san.getDuplicate() == duplicate)
return ::testing::AssertionSuccess();
return ::testing::AssertionFailure() << "tally is " << san.get() << ", expected good:" << good
<< " bad:" << bad << " dupe:" << duplicate;
}
class SHAMapSyncTest : public ::testing::Test
{
protected:
@@ -81,8 +133,517 @@ protected:
return true;
}
/**
* A sync filter that records every node it is told about, and serves back
* only the ones it was explicitly asked to hold.
*
* Serving is opt-in: the sync path consults the filter before deciding
* a node is missing, so each case serves only the node it wants resolved.
*/
class RecordingFilter : public SHAMapSyncFilter
{
public:
// What one gotNode() call was told, in the order the calls arrived.
struct Report
{
bool fromFilter;
SHAMapHash hash;
std::uint32_t ledgerSeq;
};
void
gotNode(
bool fromFilter,
SHAMapHash const& hash,
std::uint32_t ledgerSeq,
Blob&&, // NOLINT(cppcoreguidelines-rvalue-reference-param-not-moved)
SHAMapNodeType) const override
{
reports_.push_back({.fromFilter = fromFilter, .hash = hash, .ledgerSeq = ledgerSeq});
}
[[nodiscard]] std::optional<Blob>
getNode(SHAMapHash const& hash) const override
{
if (auto const it = served_.find(hash); it != served_.end())
return it->second;
return std::nullopt;
}
/**
* Offer a node back to the map, as a fetch pack does.
*
* @param node The node to serve, keyed by its own hash.
*/
void
serve(SHAMapTreeNodePtr const& node)
{
Serializer s;
node->serializeWithPrefix(s);
served_.emplace(node->getHash(), s.modData());
}
[[nodiscard]] std::vector<Report> const&
reports() const
{
return reports_;
}
private:
// Mutable because the whole interface is const: a filter is handed to the map by
// const pointer, so recording has to happen through one.
mutable std::vector<Report> reports_;
std::map<SHAMapHash, Blob> served_;
};
/**
* A root inner node with all 16 branches occupied and not one of them
* resolvable.
*
* A walk of a backed map posts an asynchronous read for every branch in a
* single pass, so the nodestore reader threads run finishFetch() for the
* same map at the same time.
*/
struct WideRoot
{
SHAMapTreeNodePtr node;
SHAMapHash hash;
WideRoot()
{
std::vector<InnerChild> children;
children.reserve(SHAMap::kBranchFactor);
for (auto branch = 0u; branch < SHAMap::kBranchFactor; ++branch)
{
// Derived from the branch, so each posts its own read and no hash collides with a
// real node's.
UInt256 childHash;
childHash.begin()[0] = 0xFA;
childHash.begin()[1] = 0xB1;
childHash.begin()[2] = static_cast<unsigned char>(branch);
children.push_back({.branch = branch, .hash = SHAMapHash{childHash}});
}
node = makeFullInnerNode(children);
hash = node->getHash();
}
};
};
// Only a leaf may sit at kLeafDepth. An inner node there is reported as bad data and leaves the
// map invalid.
TEST_F(SHAMapSyncTest, inner_node_at_leaf_depth)
{
TestNodeFamily f{j_};
DeepChain const chain;
SHAMap map{SHAMapType::FREE, f};
map.setSynching();
ASSERT_TRUE(chain.fill(map));
ASSERT_TRUE(map.isValid());
auto const result = chain.addOffendingNode(map);
EXPECT_TRUE(tallyIs(result, 0, 1, 0));
EXPECT_FALSE(result.isGood());
EXPECT_FALSE(map.isValid());
// Invalid is terminal, so setImmutable() refuses.
EXPECT_FALSE(map.setImmutable());
}
// A node the descent cannot hook in is bad data: the batch counts it bad and the map stays usable
// for another sender. All three refusals are covered: a depth the node does not sit at, a branch
// the root leaves empty, and a hash the root does not name.
TEST_F(SHAMapSyncTest, node_that_cannot_be_hooked_is_bad_data)
{
TestNodeFamily f{j_};
DeepChain const chain;
SHAMap map{SHAMapType::FREE, f};
map.setSynching();
ASSERT_TRUE(map.addRootNode(chain.rootHash, chain.nodeAt(0), nullptr).isGood());
// nodeAt(1) is the node the root is missing and its hash matches, but we claim depth 2.
auto const wrongDepth = map.addKnownNode(SHAMapNodeID{2, UInt256{}}, chain.nodeAt(1), nullptr);
EXPECT_TRUE(tallyIs(wrongDepth, 0, 1, 0));
EXPECT_FALSE(wrongDepth.isUseful());
// The chain sits on branch 0 at every depth, so a node claiming a position on branch 1 asks the
// descent to follow a branch the root leaves empty.
UInt256 otherBranch;
otherBranch.begin()[0] = 0x10;
auto const emptyBranch =
map.addKnownNode(SHAMapNodeID{1, otherBranch}, chain.nodeAt(1), nullptr);
EXPECT_TRUE(tallyIs(emptyBranch, 0, 1, 0));
// The right position this time, but the data hashes to something other than the child the root
// says belongs there.
auto const corrupt = map.addKnownNode(SHAMapNodeID{1, UInt256{}}, chain.nodeAt(2), nullptr);
EXPECT_TRUE(tallyIs(corrupt, 0, 1, 0));
// The verdict is bad data alone, so the map stays usable.
EXPECT_TRUE(map.isValid());
}
// A root is installed once and a map is synced against one hash, so a root offered under a hash the
// map does not hold names another tree and is bad data. The same root under the hash the map does
// hold is the duplicate it is.
TEST_F(SHAMapSyncTest, add_root_node_judges_the_hash_asked_for)
{
TestNodeFamily f{j_};
DeepChain const chain;
SHAMap map{SHAMapType::FREE, f};
map.setSynching();
ASSERT_TRUE(map.addRootNode(chain.rootHash, chain.nodeAt(0), nullptr).isGood());
// The same root under the hash the map holds: already held, and reported as such.
auto const same = map.addRootNode(chain.rootHash, chain.nodeAt(0), nullptr);
EXPECT_TRUE(tallyIs(same, 0, 0, 1));
EXPECT_TRUE(same.isGood());
// A hash the map does not hold. nodeAt(1) is a real node of the same chain, so this is a
// well-formed hash that names another tree.
auto const other = map.addRootNode(chain.nodeAt(1)->getHash(), chain.nodeAt(0), nullptr);
EXPECT_TRUE(tallyIs(other, 0, 1, 0));
EXPECT_FALSE(other.isGood());
EXPECT_TRUE(other.isInvalid());
// The refusal is about the hash asked for, so the root the map holds stays in place and the map
// stays usable.
EXPECT_EQ(map.getHash(), chain.rootHash);
EXPECT_TRUE(map.isValid());
}
// The verdict outranks the full-below cache. That cache is keyed by node hash and shared by every
// map of a family, and a hash covers a node's children but not its depth, so an earlier walk can
// mark the same subtree hash complete at one depth while this map reaches it at kLeafDepth, with no
// collision involved. The descent therefore skips the lookup at that boundary and reaches the depth
// verdict first. This case seeds the entry a lookup would match, so dropping the skip turns the
// verdict back into a duplicate.
TEST_F(SHAMapSyncTest, map_invalidating_node_is_judged_before_the_cache_is_read)
{
TestNodeFamily f{j_};
DeepChain const chain;
SHAMap map{SHAMapType::FREE, f};
map.setSynching();
ASSERT_TRUE(chain.fill(map));
ASSERT_TRUE(map.isValid());
// This case seeds the entry a lookup at the boundary would match: the offending node's own
// hash. The descent skips the lookup there, so the entry is never read and the depth verdict
// stands. Drop the skip and the hit returns for the whole branch, so the tally below becomes
// a duplicate.
f.getFullBelowCache()->insert(chain.nodeAt(SHAMap::kLeafDepth)->getHash().asUInt256());
auto const result = chain.addOffendingNode(map);
EXPECT_TRUE(tallyIs(result, 0, 1, 0));
EXPECT_FALSE(result.isGood());
EXPECT_FALSE(map.isValid());
EXPECT_FALSE(map.setImmutable());
}
// An invalid tx map must also stop the enclosing ledger from being marked immutable, since an
// immutable ledger is treated as persistable.
TEST_F(SHAMapSyncTest, invalid_tx_map_blocks_immutable_ledger)
{
TestNodeFamily f{j_};
DeepChain const chain;
Ledger ledger{1, NetClock::time_point{}, noAmendments(), Fees{}, f};
ASSERT_FALSE(ledger.isImmutable());
ledger.txMap().setSynching();
ASSERT_TRUE(chain.fill(ledger.txMap()));
auto const result = chain.addOffendingNode(ledger.txMap());
ASSERT_TRUE(tallyIs(result, 0, 1, 0));
ASSERT_FALSE(ledger.txMap().isValid());
// The state map is untouched, so only the transaction map can be refusing.
ASSERT_TRUE(ledger.stateMap().isValid());
EXPECT_FALSE(ledger.setImmutable());
EXPECT_FALSE(ledger.isImmutable());
}
// The same for the state map, which is the second operand of the one expression
// Ledger::setImmutable() tests both maps in.
TEST_F(SHAMapSyncTest, invalid_state_map_blocks_immutable_ledger)
{
TestNodeFamily f{j_};
DeepChain const chain;
Ledger ledger{1, NetClock::time_point{}, noAmendments(), Fees{}, f};
ASSERT_FALSE(ledger.isImmutable());
ledger.stateMap().setSynching();
ASSERT_TRUE(chain.fill(ledger.stateMap()));
auto const result = chain.addOffendingNode(ledger.stateMap());
ASSERT_TRUE(tallyIs(result, 0, 1, 0));
ASSERT_FALSE(ledger.stateMap().isValid());
// The transaction map is untouched, so only the state map can be refusing.
ASSERT_TRUE(ledger.txMap().isValid());
EXPECT_FALSE(ledger.setImmutable());
EXPECT_FALSE(ledger.isImmutable());
}
// A refusal leaves the header exactly as it was. setImmutable() derives the map hashes from the
// maps and then the ledger hash from the header, and writes them only after every check has
// passed. The up-front check covers this case, and the re-test after the maps are settled shares
// the rule, which is why the header is written only once that one has passed too.
TEST_F(SHAMapSyncTest, refused_settle_leaves_the_header_alone)
{
TestNodeFamily f{j_};
DeepChain const chain;
// Not the header constructor: this one derives its map hashes, which is what must not happen.
Ledger ledger{1, NetClock::time_point{}, noAmendments(), Fees{}, f};
ASSERT_FALSE(ledger.isImmutable());
ASSERT_TRUE(ledger.header().txHash.isZero());
ASSERT_TRUE(ledger.header().accountHash.isZero());
auto const hashBefore = ledger.header().hash;
// A transaction map that hashes to something, so a derived header hash differs from the one
// the ledger has now.
ASSERT_TRUE(ledger.txMap().addItem(SHAMapNodeType::TnTransactionNm, makeRandomAS()));
ASSERT_TRUE(ledger.txMap().getHash().isNonZero());
// And a state map the chain abandons, so settling has to refuse.
ledger.stateMap().setSynching();
ASSERT_TRUE(chain.fill(ledger.stateMap()));
ASSERT_TRUE(chain.addOffendingNode(ledger.stateMap()).isInvalid());
ASSERT_FALSE(ledger.stateMap().isValid());
EXPECT_FALSE(ledger.setImmutable());
// The map hashes, the ledger hash and the flag all still hold their pre-refusal values.
EXPECT_FALSE(ledger.isImmutable());
EXPECT_TRUE(ledger.header().txHash.isZero());
EXPECT_TRUE(ledger.header().accountHash.isZero());
EXPECT_EQ(ledger.header().hash, hashBefore);
}
// Invalid is terminal: setImmutable() and clearSynching() offer no way back out of it, however
// many times they are called. setSynching() is left alone, since it is unreachable on an invalid
// map today and says so with an UNREACHABLE.
TEST_F(SHAMapSyncTest, invalid_state_is_terminal)
{
TestNodeFamily f{j_};
DeepChain const chain;
SHAMap map{SHAMapType::FREE, f};
map.setSynching();
ASSERT_TRUE(chain.fill(map));
ASSERT_TRUE(chain.addOffendingNode(map).isInvalid());
ASSERT_FALSE(map.isValid());
// Repeated attempts must each fail, and must not leave the map reporting a valid state.
for (auto attempt = 0; attempt < 3; ++attempt)
{
EXPECT_FALSE(map.setImmutable()) << "attempt " << attempt;
EXPECT_FALSE(map.isValid()) << "attempt " << attempt;
}
// Nor does clearSynching(), which keeps an abandoned map from being moved back to Modifying and
// passing isValid() again. It refuses rather than treating that as unreachable, since a
// concurrent walk can invalidate a map between a caller's own check and this call.
for (auto attempt = 0; attempt < 3; ++attempt)
{
map.clearSynching();
EXPECT_FALSE(map.isValid()) << "attempt " << attempt;
}
// isSynching() answers false for an invalid map.
EXPECT_FALSE(map.isSynching());
}
// Immutable holds against clearSynching(): a walk that ends after the ledger settled leaves the map
// settled rather than reopening it for modification.
TEST_F(SHAMapSyncTest, clear_synching_leaves_an_immutable_map_immutable)
{
TestNodeFamily f{j_};
SHAMap map{SHAMapType::FREE, f};
map.setSynching();
ASSERT_TRUE(map.setImmutable());
ASSERT_TRUE(map.isImmutable());
map.clearSynching();
EXPECT_TRUE(map.isImmutable());
EXPECT_FALSE(map.isSynching());
}
// A snapshot shares the source map's root, so Invalid carries over to it.
TEST_F(SHAMapSyncTest, snapshot_of_invalid_map_stays_invalid)
{
TestNodeFamily f{j_};
DeepChain const chain;
SHAMap map{SHAMapType::FREE, f};
map.setSynching();
ASSERT_TRUE(chain.fill(map));
auto const result = chain.addOffendingNode(map);
ASSERT_TRUE(tallyIs(result, 0, 1, 0));
ASSERT_FALSE(map.isValid());
// Both flavors: the immutable snapshot is the one the store reads, and the
// mutable one is the one that sets Modifying.
for (bool const isMutable : {false, true})
{
auto const snapshot = map.snapShot(isMutable);
ASSERT_TRUE(snapshot != nullptr);
EXPECT_FALSE(snapshot->isValid()) << "isMutable " << isMutable;
EXPECT_FALSE(snapshot->setImmutable()) << "isMutable " << isMutable;
}
// A snapshot of a sound map is unaffected.
SHAMap valid{SHAMapType::FREE, f};
valid.addItem(SHAMapNodeType::TnAccountState, makeRandomAS());
EXPECT_TRUE(valid.snapShot(false)->isValid());
EXPECT_TRUE(valid.snapShot(true)->isValid());
}
// A map marked complete in the database withdraws that claim the first time a read misses, and
// reports the miss once so the ledger can be re-acquired. Sixteen unresolvable branches are posted
// in one pass, so with four reader threads the misses overlap and finishFetch() runs concurrently
// for a single map.
//
// This pins the report path: a miss withdraws the flag and reports once. The nightly
// ThreadSanitizer job that PR 8245 adds covers the ordering.
TEST_F(SHAMapSyncTest, full_flag_is_withdrawn_once_by_concurrent_readers)
{
static constexpr auto kRounds = 8uz;
static constexpr auto kReadThreads = 4;
for (auto round = 0uz; round < kRounds; ++round)
{
TestNodeFamily f{j_, kReadThreads};
WideRoot const root;
// Backed, so descendAsync() posts real asynchronous reads rather than resolving inline.
SHAMap map{SHAMapType::FREE, f};
map.setSynching();
ASSERT_TRUE(map.addRootNode(root.hash, root.node, nullptr).isGood());
// The claim the first miss has to withdraw.
map.setFull();
// A null filter, so every branch is read from a database that lacks it.
EXPECT_EQ(map.getMissingNodes(kMaxNodesPerRequest, nullptr).size(), SHAMap::kBranchFactor)
<< "round " << round;
EXPECT_EQ(f.missingBySeqReports(), 1uz) << "round " << round;
}
}
// Every node the sync path hands to a filter carries the map's ledger sequence. All three call
// sites are covered: a root taken from a peer, a node taken from a peer, and a node the walk
// resolved out of the filter itself.
TEST_F(SHAMapSyncTest, sync_filter_is_told_the_ledger_sequence)
{
static constexpr std::uint32_t kLedgerSeq = 7;
TestNodeFamily f{j_};
// A three-level chain, built bottom-up so each hash covers the one below it. Every node keeps
// its one child on branch 0. The deepest points at a child the fixture withholds, so the walk
// always has something to ask for.
auto const deepest = makeCompressedInnerNode({{.branch = 0u, .hash = SHAMapHash{UInt256{1}}}});
auto const middle = makeCompressedInnerNode({{.branch = 0u, .hash = deepest->getHash()}});
auto const root = makeCompressedInnerNode({{.branch = 0u, .hash = middle->getHash()}});
// Unbacked, so the filter is the only source, and a node it withholds counts as missing.
SHAMap map{SHAMapType::FREE, f};
map.setUnbacked();
map.setSynching();
map.setLedgerSeq(kLedgerSeq);
RecordingFilter filter;
ASSERT_TRUE(map.addRootNode(root->getHash(), root, &filter).isGood());
ASSERT_TRUE(map.addKnownNode(SHAMapNodeID{1, UInt256{}}, middle, &filter).isUseful());
// The deepest node becomes resolvable only now, after the two above were added directly.
filter.serve(deepest);
auto const missing = map.getMissingNodes(kMaxNodesPerRequest, &filter);
// The walk resolved the deepest node through the filter and then asked for its child.
ASSERT_EQ(missing.size(), 1u);
EXPECT_EQ(missing[0].second, UInt256{1});
ASSERT_EQ(filter.reports().size(), 3u);
// The two nodes taken from a peer, which the filter is told about so it can store them.
EXPECT_FALSE(filter.reports()[0].fromFilter);
EXPECT_EQ(filter.reports()[0].hash, root->getHash());
EXPECT_FALSE(filter.reports()[1].fromFilter);
EXPECT_EQ(filter.reports()[1].hash, middle->getHash());
// The one the walk read back out of the filter, which is reported as such.
EXPECT_TRUE(filter.reports()[2].fromFilter);
EXPECT_EQ(filter.reports()[2].hash, deepest->getHash());
for (auto const& report : filter.reports())
EXPECT_EQ(report.ledgerSeq, kLedgerSeq) << "hash " << report.hash;
}
// Ledger::setFull() publishes each map's ledger sequence alongside the flag that lets the first
// nodestore miss report a gap. The sequence is what the lookup resolving that gap reads.
//
// This pins that setFull() sets the sequence. The nightly ThreadSanitizer job that PR 8245
// adds covers the order of the two stores.
TEST_F(SHAMapSyncTest, ledger_set_full_publishes_the_ledger_sequence)
{
static constexpr std::uint32_t kLedgerSeq = 7;
TestNodeFamily f{j_};
LedgerHeader header;
header.seq = kLedgerSeq;
// Non-zero, so the map has a root to look for and the lookup can miss.
header.txHash = UInt256{1};
header.hash = calculateLedgerHash(header);
Ledger ledger{header, noAmendments(), f};
// The constructor already looked for that root and missed. The map becomes complete only at
// setFull() below, so the report count is still zero here.
ASSERT_EQ(f.missingBySeqReports(), 0uz);
ledger.setFull();
// Still missing, and now the map has a claim to withdraw, so the gap is reported.
// TestNodeFamily throws in place of the real family's re-acquisition, which finishFetch()
// logs and swallows.
EXPECT_FALSE(ledger.txMap().fetchRoot(SHAMapHash{header.txHash}, nullptr));
EXPECT_EQ(f.missingBySeqReports(), 1uz);
EXPECT_EQ(f.missingBySeqRefNum(), kLedgerSeq);
}
TEST_F(SHAMapSyncTest, sync)
{
TestNodeFamily f{j_}, f2{j_};
@@ -92,7 +653,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)
{
@@ -105,7 +665,7 @@ TEST_F(SHAMapSyncTest, sync)
ASSERT_TRUE(confuseMap(source, kNodesToConfuse));
source.invariants();
source.setImmutable();
ASSERT_TRUE(source.setImmutable());
std::size_t count = 0;
source.visitLeaves([&count]([[maybe_unused]] auto const& item) { ++count; });
@@ -180,4 +740,50 @@ TEST_F(SHAMapSyncTest, sync)
destination.invariants();
}
// The duplicate verdict also answers for a node that did not need to be added, which is the half
// of incDuplicate()'s contract that "already held" does not cover. addKnownNode() reaches it when
// the map has stopped taking nodes, and returns there before looking at the offer at all.
//
// An A/B on one offer and two maps of the same shape, so the synching state is the only difference
// between the two verdicts. Neither map holds the offered node, so a count that meant "already
// held" would be wrong for it.
TEST_F(SHAMapSyncTest, add_known_node_reports_duplicate_once_a_map_stops_synching)
{
TestNodeFamily f{j_};
// A well-formed inner node with one child. Its contents do not matter: both verdicts below
// are reached without the map reading them.
auto const makeOffer = [] {
Serializer s;
s.addBitString(UInt256{1});
s.add8(0);
s.add8(kWireTypeCompressedInner);
return SHAMapTreeNode::makeFromWire(makeSlice(s.peekData()));
};
ASSERT_TRUE(makeOffer());
SHAMapNodeID const target{1, UInt256{}};
// Synching, so the offer is examined, and refused because the empty root has no branch to
// hook it onto. The point is that the map looked.
SHAMap synching{SHAMapType::FREE, f};
synching.setSynching();
auto const examined = synching.addKnownNode(target, makeOffer(), nullptr);
EXPECT_TRUE(examined.isInvalid());
EXPECT_EQ(examined.getDuplicate(), 0);
// The same offer to a map that has stopped taking nodes.
SHAMap stopped{SHAMapType::FREE, f};
stopped.setSynching();
stopped.clearSynching();
ASSERT_FALSE(stopped.isSynching());
auto const notNeeded = stopped.addKnownNode(target, makeOffer(), nullptr);
EXPECT_EQ(notNeeded.getGood(), 0);
EXPECT_EQ(notNeeded.getBad(), 0);
EXPECT_EQ(notNeeded.getDuplicate(), 1);
EXPECT_TRUE(notNeeded.isGood()) << "a node that was not needed counts on the accepted side";
}
} // namespace xrpl::tests

View File

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

View File

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

View File

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

View File

@@ -30,9 +30,9 @@
namespace xrpl {
// A ledger we are trying to acquire
class InboundLedger final : public TimeoutCounter,
public std::enable_shared_from_this<InboundLedger>,
public CountedObject<InboundLedger>
class InboundLedger : public TimeoutCounter,
public std::enable_shared_from_this<InboundLedger>,
public CountedObject<InboundLedger>
{
public:
using ClockType = beast::AbstractClock<std::chrono::steady_clock>;
@@ -44,13 +44,29 @@ public:
CONSENSUS // We believe the consensus round requires this ledger
};
/**
* How long to wait between retries, and so how long one timeout takes.
*/
static constexpr std::chrono::milliseconds kRetryInterval{3000};
/**
* @param app The application to run in.
* @param hash The ledger to acquire.
* @param seq Its sequence, or zero if not known yet.
* @param reason Why it is being acquired.
* @param clock The clock touch() records against.
* @param peerSet Which peers to ask, and how to reach them.
* @param retryInterval How long to wait between retries. TimeoutCounter
* requires more than 10ms and less than 30s.
*/
InboundLedger(
Application& app,
UInt256 const& hash,
std::uint32_t seq,
Reason reason,
ClockType&,
std::unique_ptr<PeerSet> peerSet);
ClockType& clock,
std::unique_ptr<PeerSet> peerSet,
std::chrono::milliseconds retryInterval = kRetryInterval);
~InboundLedger() override;
@@ -76,10 +92,19 @@ public:
return failed_;
}
/**
* The acquired ledger.
*
* A failed acquisition may still hold a partially built ledger, which
* getJson() reports on, so ledger_ is kept while this answers nullptr.
*
* @return The ledger, or nullptr before a header is obtained and once the
* acquisition has failed.
*/
std::shared_ptr<Ledger const>
getLedger() const
{
return ledger_;
return failed_ ? nullptr : ledger_;
}
std::uint32_t
@@ -119,15 +144,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. Call under mtx_, which the flags written here require.
*/
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 +181,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 +195,17 @@ private:
bool
takeHeader(std::string_view data);
/**
* Fail the acquisition when the header's account hash is zero. No ledger
* has an empty state map, so such a header cannot name a ledger. Both
* tryDB() and takeHeader() judge the header here.
*
* @return Whether the acquisition was failed. Then failed_ is set and
* ledger_ is null.
*/
bool
failOnZeroAccountHash();
void
receiveNode(
std::shared_ptr<Peer> const& peer,

View File

@@ -6,6 +6,7 @@
#include <xrpl/basics/Log.h>
#include <xrpl/basics/chrono.h>
#include <xrpl/basics/contract.h>
#include <xrpl/beast/utility/Journal.h>
#include <xrpl/beast/utility/instrumentation.h>
#include <xrpl/ledger/ApplyView.h>
@@ -76,7 +77,18 @@ buildLedgerImpl(
XRPL_ASSERT(
built->header().seq < kXrpLedgerEarliestFees || built->read(keylet::feeSettings()),
"xrpl::buildLedgerImpl : valid ledger fees");
built->setAccepted(closeTime, closeResolution, closeTimeCorrect);
// The invariant: a ledger this function returns has both maps immutable, which is what
// setAccepted() reports on (see Ledger::setImmutable()). Nothing downstream re-checks it.
//
// logicError() rather than UNREACHABLE(): the consensus caller acts on this ledger straight
// away, so the stop has to name this site in every build rather than only where assertions
// are on. See RCLConsensus::Adaptor::buildLCL().
if (!built->setAccepted(closeTime, closeResolution, closeTimeCorrect))
{
// LCOV_EXCL_START
logicError("buildLedgerImpl: accepted ledger map is invalid");
// LCOV_EXCL_STOP
}
return built;
}

View File

@@ -54,8 +54,6 @@
namespace xrpl {
using namespace std::chrono_literals;
static constexpr auto kPeerCountStart = 5; // Number of peers to start with
static constexpr auto kPeerCountAdd = 3; // Number of peers to add on a timeout
static constexpr auto kLedgerTimeoutRetriesMax = 6; // how many timeouts before we give up
@@ -65,20 +63,18 @@ static constexpr auto kMissingNodesFind = 256; // Number of nodes to find initi
static constexpr auto kReqNodesReply = 128; // Number of nodes to request for a reply
static constexpr auto kReqNodes = 12; // Number of nodes to request blindly
// millisecond for each ledger timeout
constexpr auto kLedgerAcquireTimeout = 3000ms;
InboundLedger::InboundLedger(
Application& app,
UInt256 const& hash,
std::uint32_t seq,
Reason reason,
ClockType& clock,
std::unique_ptr<PeerSet> peerSet)
std::unique_ptr<PeerSet> peerSet,
std::chrono::milliseconds retryInterval)
: TimeoutCounter(
app,
hash,
kLedgerAcquireTimeout,
retryInterval,
{.jobType = JtLedgerData, .jobName = "InboundLedger", .jobLimit = 5},
app.getJournal("InboundLedger"))
, clock_(clock)
@@ -97,8 +93,12 @@ InboundLedger::init(ScopedLockType& collectionLock)
collectionLock.unlock();
tryDB(app_.getNodeFamily().db());
// done() is what wakes whatever is waiting and records the hash in recentFailures_.
if (failed_)
{
done();
return;
}
if (!complete_)
{
@@ -112,7 +112,21 @@ InboundLedger::init(ScopedLockType& collectionLock)
XRPL_ASSERT(
ledger_->header().seq < kXrpLedgerEarliestFees || ledger_->read(keylet::feeSettings()),
"xrpl::InboundLedger::init : valid ledger fees");
ledger_->setImmutable();
// tryDB() verified both maps before setting complete_ and mtx_ has been held since, so
// nothing can have invalidated them.
if (!ledger_->setImmutable())
{
// LCOV_EXCL_START
// Recorded before the UNREACHABLE, which is not guaranteed to stop here.
// Withdrawn as well as failed, to match done() and TransactionAcquire::done(),
// for a caller that checks complete_ before failed_.
complete_ = false;
failed_ = true;
done();
UNREACHABLE("xrpl::InboundLedger::init : map is invalid");
return;
// LCOV_EXCL_STOP
}
if (reason_ == Reason::HISTORY)
return;
@@ -241,7 +255,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 +328,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))
{
@@ -328,12 +341,20 @@ InboundLedger::tryDB(node_store::Database& srcDB)
if (haveTransactions_ && haveState_)
{
JLOG(journal_.debug()) << "Had everything locally";
complete_ = true;
XRPL_ASSERT(
ledger_->header().seq < kXrpLedgerEarliestFees || ledger_->read(keylet::feeSettings()),
"xrpl::InboundLedger::tryDB : valid ledger fees");
ledger_->setImmutable();
// Settled before complete_ is published, so a caller that reads the flag never sees a
// ledger this function has not finished with.
if (!ledger_->setImmutable())
{
JLOG(journal_.warn()) << "Ledger " << hash_ << " found locally is invalid";
failed_ = true;
return;
}
JLOG(journal_.debug()) << "Had everything locally";
complete_ = true;
}
}
@@ -419,6 +440,40 @@ InboundLedger::done()
signaled_ = true;
touch();
// Settled before the outcome below is logged or acted on, so nothing reports a ledger this
// function has since refused.
if (complete_ && !failed_ && ledger_)
{
XRPL_ASSERT(
ledger_->header().seq < kXrpLedgerEarliestFees || ledger_->read(keylet::feeSettings()),
"xrpl::InboundLedger::done : valid ledger fees");
// Recovers rather than asserting: peer data produces this verdict, so a caller cannot know
// its map is still sound. Recovering rather than asserting is therefore the contract.
// Best-effort even so: setInvalid() outranks Immutable, so a walk that reaches the verdict
// after both maps have been settled leaves an immutable ledger with an invalid map. It
// narrows the window rather than closing it.
if (!ledger_->setImmutable())
{
JLOG(journal_.warn()) << "Acquired ledger " << hash_ << " is invalid";
// Withdrawn as well as failed, so a caller that already read complete_ - or that checks
// it before failed_ - cannot go on treating this ledger as delivered.
complete_ = false;
failed_ = true;
}
else
{
switch (reason_)
{
case Reason::HISTORY:
app_.getInboundLedgers().onLedgerFetched();
break;
default:
app_.getLedgerMaster().storeLedger(ledger_);
break;
}
}
}
JLOG(journal_.debug()) << "Acquire " << hash_ << (failed_ ? " fail " : " ")
<< ((timeouts_ == 0)
? std::string()
@@ -427,24 +482,7 @@ InboundLedger::done()
XRPL_ASSERT(complete_ || failed_, "xrpl::InboundLedger::done : complete or failed");
if (complete_ && !failed_ && ledger_)
{
XRPL_ASSERT(
ledger_->header().seq < kXrpLedgerEarliestFees || ledger_->read(keylet::feeSettings()),
"xrpl::InboundLedger::done : valid ledger fees");
ledger_->setImmutable();
switch (reason_)
{
case Reason::HISTORY:
app_.getInboundLedgers().onLedgerFetched();
break;
default:
app_.getLedgerMaster().storeLedger(ledger_);
break;
}
}
// We hold the PeerSet lock, so must dispatch
// mtx_ is held, so this may only post the work rather than do it.
app_.getJobQueue().addJob(JtLedgerData, "AcqDone", [self = shared_from_this()]() {
if (self->complete_ && !self->failed_)
{
@@ -497,6 +535,7 @@ InboundLedger::trigger(std::shared_ptr<Peer> const& peer, TriggerReason reason)
if (failed_)
{
JLOG(journal_.warn()) << " failed local for " << hash_;
done();
return;
}
}
@@ -730,7 +769,8 @@ InboundLedger::trigger(std::shared_ptr<Peer> const& peer, TriggerReason reason)
{
JLOG(journal_.debug()) << "Done:" << (complete_ ? " complete" : "")
<< (failed_ ? " failed " : " ") << ledger_->header().seq;
sl.unlock();
// Called with mtx_ still held, so the flags done() writes are not written unlocked; mtx_
// is recursive, so a caller that already holds it further up is unaffected.
done();
}
}
@@ -774,15 +814,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 +854,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 +874,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 +1158,13 @@ InboundLedger::processData(std::shared_ptr<Peer> peer, protocol::TMLedgerData co
return -1;
}
// takeHeader() failed the acquisition. Nothing else signals that after isDone().
if (failed_)
{
done();
return 0;
}
san.incUseful();
}

View File

@@ -115,8 +115,15 @@ loadLedgerHelper(
return ledger;
}
/**
* Settle a ledger just loaded from local storage, or discard it.
*
* @param ledger The ledger to settle. Cleared on failure, so the caller cannot
* hand out one that is still mutable.
* @param j Where to log a refusal.
*/
static void
finishLoadByIndexOrHash(std::shared_ptr<Ledger> const& ledger, beast::Journal j)
finishLoadByIndexOrHash(std::shared_ptr<Ledger>& ledger, beast::Journal j)
{
if (!ledger)
return;
@@ -124,7 +131,18 @@ finishLoadByIndexOrHash(std::shared_ptr<Ledger> const& ledger, beast::Journal j)
XRPL_ASSERT(
ledger->header().seq < kXrpLedgerEarliestFees || ledger->read(keylet::feeSettings()),
"xrpl::finishLoadByIndexOrHash : valid ledger fees");
ledger->setImmutable();
// Loaded locally. See Ledger::setImmutable().
if (!ledger->setImmutable())
{
// LCOV_EXCL_START
JLOG(j.error()) << "Invalid map for ledger " << ledger->header().seq
<< "; not marking it as loaded";
UNREACHABLE("xrpl::finishLoadByIndexOrHash : map is invalid");
// Discarded rather than left un-full, since nothing gates usability on the full flag.
ledger.reset();
return;
// LCOV_EXCL_STOP
}
JLOG(j.trace()) << "Loaded ledger: " << to_string(ledger->header().hash);

View File

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

View File

@@ -8,6 +8,7 @@
#include <xrpl/basics/Log.h>
#include <xrpl/basics/base_uint.h>
#include <xrpl/beast/utility/instrumentation.h>
#include <xrpl/core/Job.h>
#include <xrpl/server/NetworkOPs.h>
#include <xrpl/shamap/SHAMap.h>
@@ -18,6 +19,7 @@
#include <xrpl.pb.h>
#include <algorithm>
#include <chrono>
#include <cstddef>
#include <exception>
#include <memory>
@@ -26,22 +28,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))
@@ -53,16 +51,29 @@ TransactionAcquire::TransactionAcquire(
void
TransactionAcquire::done()
{
// We hold a PeerSet lock and so cannot do real work here
// mtx_ is held, so this may only post real work rather than do it.
if (failed_)
{
JLOG(journal_.debug()) << "Failed to acquire TX set " << hash_;
}
else if (!map_->setImmutable())
{
// trigger() verified the map before setting complete_ and mtx_ has been held since, and
// nothing walks this map with the lock released, so it is still valid. UNREACHABLE for
// that reason, with the flags still withdrawn, since UNREACHABLE need not stop here.
// LCOV_EXCL_START
// Withdraw complete_ alongside the failure, for trigger() and takeNodes(), which both
// check complete_ before failed_.
complete_ = false;
failed_ = true;
JLOG(journal_.debug()) << "Failed to acquire TX set " << hash_;
UNREACHABLE("xrpl::TransactionAcquire::done : map is invalid");
// LCOV_EXCL_STOP
}
else
{
JLOG(journal_.debug()) << "Acquired TX set " << hash_;
map_->setImmutable();
UInt256 const& hash(hash_);
std::shared_ptr<SHAMap> const& map(map_);
@@ -79,7 +90,7 @@ TransactionAcquire::done()
}
void
TransactionAcquire::onTimer(bool progress, ScopedLockType& psl)
TransactionAcquire::onTimer(bool progress, ScopedLockType&)
{
if (timeouts_ > kMaxTimeouts)
{

View File

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

View File

@@ -1702,7 +1702,14 @@ ApplicationImp::startGenesisLedger()
XRPL_ASSERT(
next->header().seq < kXrpLedgerEarliestFees || next->read(keylet::feeSettings()),
"xrpl::ApplicationImp::startGenesisLedger : valid ledger fees");
next->setImmutable();
// Built locally. See Ledger::setImmutable(). Failed here rather than at storeLedger() below,
// which would name the wrong site.
if (!next->setImmutable())
{
// LCOV_EXCL_START
logicError("startGenesisLedger: genesis ledger map is invalid");
// LCOV_EXCL_STOP
}
openLedger_.emplace(next, cachedSLEs_, logs_->journal("OpenLedger"));
ledgerMaster_->storeLedger(next);
ledgerMaster_->switchLCL(next);
@@ -1724,7 +1731,16 @@ ApplicationImp::getLastFullLedger()
XRPL_ASSERT(
ledger->header().seq < kXrpLedgerEarliestFees || ledger->read(keylet::feeSettings()),
"xrpl::ApplicationImp::getLastFullLedger : valid ledger fees");
ledger->setImmutable();
// Loaded locally. See Ledger::setImmutable().
if (!ledger->setImmutable())
{
// LCOV_EXCL_START
JLOG(j.error()) << "Last full ledger " << seq << " has an invalid map; ignoring it";
UNREACHABLE("xrpl::ApplicationImp::getLastFullLedger : map is invalid");
// Must not fall through: that would mark a damaged ledger validated.
return {};
// LCOV_EXCL_STOP
}
if (getLedgerMaster().haveLedger(seq))
ledger->setValidated();
@@ -1876,7 +1892,16 @@ ApplicationImp::loadLedgerFromFile(std::string const& name)
loadLedger->header().seq < kXrpLedgerEarliestFees ||
loadLedger->read(keylet::feeSettings()),
"xrpl::ApplicationImp::loadLedgerFromFile : valid ledger fees");
loadLedger->setAccepted(closeTime, closeTimeResolution, !closeTimeEstimated);
// Built locally. See Ledger::setImmutable(). The caller handles a failure return, so this
// returns rather than treating it as unreachable.
if (!loadLedger->setAccepted(closeTime, closeTimeResolution, !closeTimeEstimated))
{
// LCOV_EXCL_START
JLOG(journal_.fatal()) << "Ledger from file has an invalid map";
UNREACHABLE("xrpl::ApplicationImp::loadLedgerFromFile : map is invalid");
return nullptr;
// LCOV_EXCL_STOP
}
return loadLedger;
}