mirror of
https://github.com/XRPLF/rippled.git
synced 2026-10-11 06:08:02 +00:00
Compare commits
17 Commits
mvadari/re
...
bthomee/sh
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
749cc6a959 | ||
|
|
50afd928bd | ||
|
|
b28936342a | ||
|
|
e820f5a8bf | ||
|
|
07a6f129cf | ||
|
|
db65266943 | ||
|
|
a86a2651ec | ||
|
|
08b8f62aa2 | ||
|
|
057b5788fb | ||
|
|
88ad25eb0e | ||
|
|
ccd8838d09 | ||
|
|
dcd44fa0ae | ||
|
|
a0fd44cc7b | ||
|
|
516ed4d433 | ||
|
|
a790fb3527 | ||
|
|
0514203d03 | ||
|
|
db21f68086 |
@@ -103,6 +103,7 @@ words:
|
||||
- dearmor
|
||||
- decryptor
|
||||
- dedented
|
||||
- dedup
|
||||
- deleteme
|
||||
- demultiplexer
|
||||
- deserializaton
|
||||
@@ -344,6 +345,7 @@ words:
|
||||
- unambiguity
|
||||
- unauthorizes
|
||||
- unauthorizing
|
||||
- undeserializable
|
||||
- unergonomic
|
||||
- unfetched
|
||||
- unfindable
|
||||
|
||||
@@ -62,6 +62,7 @@ libxrpl.tx > xrpl.protocol
|
||||
libxrpl.tx > xrpl.server
|
||||
libxrpl.tx > xrpl.tx
|
||||
test.app > test.jtx
|
||||
test.app > tests.libxrpl
|
||||
test.app > test.unit_test
|
||||
test.app > xrpl.basics
|
||||
test.app > xrpl.config
|
||||
|
||||
@@ -77,19 +77,24 @@ if(is_clang)
|
||||
message(STATUS " Ignorelist: ${ignorelist_path}")
|
||||
endif()
|
||||
|
||||
# Define SANITIZERS macro for BuildInfo.cpp
|
||||
# Define the SANITIZERS macro for BuildInfo.cpp, plus one of XRPL_ASAN,
|
||||
# XRPL_TSAN and XRPL_UBSAN per active sanitizer, so that code can test for a
|
||||
# specific one with #ifdef instead of parsing the dot-joined SANITIZERS string.
|
||||
set(sanitizers_list)
|
||||
if(SANITIZERS MATCHES "address")
|
||||
set(enable_asan ON)
|
||||
list(APPEND sanitizers_list "ASAN")
|
||||
target_compile_definitions(common INTERFACE XRPL_ASAN)
|
||||
endif()
|
||||
if(SANITIZERS MATCHES "thread")
|
||||
set(enable_tsan ON)
|
||||
list(APPEND sanitizers_list "TSAN")
|
||||
target_compile_definitions(common INTERFACE XRPL_TSAN)
|
||||
endif()
|
||||
if(SANITIZERS MATCHES "undefinedbehavior")
|
||||
set(enable_ubsan ON)
|
||||
list(APPEND sanitizers_list "UBSAN")
|
||||
target_compile_definitions(common INTERFACE XRPL_UBSAN)
|
||||
endif()
|
||||
|
||||
if(sanitizers_list)
|
||||
|
||||
@@ -105,6 +105,16 @@ public:
|
||||
std::vector<UInt256> const& amendments,
|
||||
Family& family);
|
||||
|
||||
/**
|
||||
* Create a ledger from a header whose maps are filled in afterwards.
|
||||
*
|
||||
* Both maps start Synching against the hashes the header carries, which are
|
||||
* input rather than derived, so setImmutable() leaves them alone.
|
||||
*
|
||||
* @param info The header to build from.
|
||||
* @param rules The rules in force.
|
||||
* @param family The SHAMap family the maps belong to.
|
||||
*/
|
||||
Ledger(LedgerHeader const& info, Rules rules, Family& family);
|
||||
|
||||
/**
|
||||
@@ -257,38 +267,79 @@ 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.
|
||||
*
|
||||
* Only a map syncing against hashes from outside can be found unsound (see
|
||||
* SHAMap::addKnownNode), so a caller whose ledger was built or loaded
|
||||
* locally treats a false return as a broken internal invariant.
|
||||
*
|
||||
* @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 and only if the header did not
|
||||
* supply them.
|
||||
* @return false if either map is Invalid, leaving the immutable flag as it
|
||||
* was and the header untouched. The ledger must then be discarded
|
||||
* rather than retried, since one map may already be immutable.
|
||||
*/
|
||||
[[nodiscard]] bool
|
||||
setImmutable(bool rehash = true);
|
||||
|
||||
/**
|
||||
* @return Whether both maps have been settled, so the ledger can no longer
|
||||
* change.
|
||||
*/
|
||||
bool
|
||||
isImmutable() const
|
||||
{
|
||||
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,8 +469,33 @@ 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_;
|
||||
|
||||
/**
|
||||
* Whether the header supplied the transaction and account hashes, so they
|
||||
* must not be derived from the maps.
|
||||
*
|
||||
* True only for a ledger built from a header, which stays verified against
|
||||
* the hash it was asked for. Fixed at construction, unlike immutable_.
|
||||
*/
|
||||
bool const mapHashesFromHeader_ = false;
|
||||
|
||||
// A SHAMap containing the transactions associated with this ledger.
|
||||
SHAMap mutable txMap_;
|
||||
|
||||
|
||||
@@ -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,42 @@ 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 setInvalid() and
|
||||
* clearSynching(), while whatever drives the acquisition reads it.
|
||||
* The walk runs with the acquisition's lock released, so this is atomic
|
||||
* rather than guarded.
|
||||
*/
|
||||
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 +180,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 +225,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;
|
||||
|
||||
@@ -296,9 +338,13 @@ public:
|
||||
* concurrency, to discover nodes referenced in the
|
||||
* SHAMap but not available locally.
|
||||
*
|
||||
* Only a leaf may occupy a position at or beyond kLeafDepth. A map that
|
||||
* breaks that is marked Invalid and the traversal is abandoned, so callers
|
||||
* ask isValid() to tell an empty result from a satisfied map.
|
||||
*
|
||||
* @param maxNodes The maximum number of found nodes to return
|
||||
* @param filter The filter to use when retrieving nodes
|
||||
* @param return The nodes known to be missing
|
||||
* @return The nodes known to be missing, or empty if the map is Invalid
|
||||
*/
|
||||
std::vector<std::pair<SHAMapNodeID, UInt256>>
|
||||
getMissingNodes(int maxNodes, SHAMapSyncFilter const* filter);
|
||||
@@ -341,6 +387,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 +409,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 +428,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 +618,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 +880,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
|
||||
|
||||
@@ -14,35 +14,112 @@ private:
|
||||
|
||||
public:
|
||||
SHAMapAddNode();
|
||||
|
||||
/**
|
||||
* Record one node that was rejected.
|
||||
*/
|
||||
void
|
||||
incInvalid();
|
||||
|
||||
/**
|
||||
* Record one node that produced a good result.
|
||||
*/
|
||||
void
|
||||
incUseful();
|
||||
|
||||
/**
|
||||
* Record one node that was not needed: the map already held it, or the map
|
||||
* or its acquisition was no longer taking nodes. isGood() counts it on the
|
||||
* accepted side.
|
||||
*/
|
||||
void
|
||||
incDuplicate();
|
||||
void
|
||||
reset();
|
||||
|
||||
/**
|
||||
* @return How many nodes produced a good result.
|
||||
*/
|
||||
[[nodiscard]] int
|
||||
getGood() const;
|
||||
|
||||
/**
|
||||
* @return How many nodes this code rejected, which is not how many the
|
||||
* batch was given: a batch that stops at its first bad node counts
|
||||
* one, and a batch that carries on counts each.
|
||||
*/
|
||||
[[nodiscard]] int
|
||||
getBad() const;
|
||||
|
||||
/**
|
||||
* @return How many nodes were not needed, whether already held or offered
|
||||
* to a map no longer taking nodes, tallied apart from the good and
|
||||
* bad counts.
|
||||
*/
|
||||
[[nodiscard]] int
|
||||
getDuplicate() const;
|
||||
|
||||
/**
|
||||
* Whether nodes that produced a good result or were not needed outnumber
|
||||
* the ones rejected.
|
||||
*
|
||||
* @return Whether the batch was good.
|
||||
*/
|
||||
[[nodiscard]] bool
|
||||
isGood() const;
|
||||
|
||||
/**
|
||||
* @return Whether at least one node in the batch was rejected.
|
||||
*/
|
||||
[[nodiscard]] bool
|
||||
isInvalid() const;
|
||||
|
||||
/**
|
||||
* @return Whether at least one node in the batch produced a good result.
|
||||
*/
|
||||
[[nodiscard]] bool
|
||||
isUseful() const;
|
||||
|
||||
/**
|
||||
* @return A verdict recording one node that was not needed.
|
||||
*/
|
||||
static SHAMapAddNode
|
||||
duplicate();
|
||||
|
||||
/**
|
||||
* @return A verdict recording one useful node.
|
||||
*/
|
||||
static SHAMapAddNode
|
||||
useful();
|
||||
|
||||
/**
|
||||
* @return A verdict recording one invalid node.
|
||||
*/
|
||||
static SHAMapAddNode
|
||||
invalid();
|
||||
|
||||
/**
|
||||
* Clear every count back to zero.
|
||||
*/
|
||||
void
|
||||
reset();
|
||||
|
||||
/**
|
||||
* Render the tally as a log line.
|
||||
*
|
||||
* @return The tally, e.g. "good:2 bad:1 dupe:1", or "no nodes processed" if
|
||||
* every count is zero.
|
||||
*/
|
||||
[[nodiscard]] std::string
|
||||
get() const;
|
||||
|
||||
/**
|
||||
* Add another verdict's counts into this one.
|
||||
*
|
||||
* @param n The verdict to add.
|
||||
* @return This verdict, updated.
|
||||
*/
|
||||
SHAMapAddNode&
|
||||
operator+=(SHAMapAddNode const& n);
|
||||
|
||||
static SHAMapAddNode
|
||||
duplicate();
|
||||
static SHAMapAddNode
|
||||
useful();
|
||||
static SHAMapAddNode
|
||||
invalid();
|
||||
|
||||
private:
|
||||
SHAMapAddNode(int good, int bad, int duplicate);
|
||||
};
|
||||
@@ -74,18 +151,30 @@ SHAMapAddNode::incDuplicate()
|
||||
++duplicate_;
|
||||
}
|
||||
|
||||
inline void
|
||||
SHAMapAddNode::reset()
|
||||
{
|
||||
good_ = bad_ = duplicate_ = 0;
|
||||
}
|
||||
|
||||
inline int
|
||||
SHAMapAddNode::getGood() const
|
||||
{
|
||||
return good_;
|
||||
}
|
||||
|
||||
inline int
|
||||
SHAMapAddNode::getBad() const
|
||||
{
|
||||
return bad_;
|
||||
}
|
||||
|
||||
inline int
|
||||
SHAMapAddNode::getDuplicate() const
|
||||
{
|
||||
return duplicate_;
|
||||
}
|
||||
|
||||
inline bool
|
||||
SHAMapAddNode::isGood() const
|
||||
{
|
||||
return (good_ + duplicate_) > bad_;
|
||||
}
|
||||
|
||||
inline bool
|
||||
SHAMapAddNode::isInvalid() const
|
||||
{
|
||||
@@ -98,22 +187,6 @@ SHAMapAddNode::isUseful() const
|
||||
return good_ > 0;
|
||||
}
|
||||
|
||||
inline SHAMapAddNode&
|
||||
SHAMapAddNode::operator+=(SHAMapAddNode const& n)
|
||||
{
|
||||
good_ += n.good_;
|
||||
bad_ += n.bad_;
|
||||
duplicate_ += n.duplicate_;
|
||||
|
||||
return *this;
|
||||
}
|
||||
|
||||
inline bool
|
||||
SHAMapAddNode::isGood() const
|
||||
{
|
||||
return (good_ + duplicate_) > bad_;
|
||||
}
|
||||
|
||||
inline SHAMapAddNode
|
||||
SHAMapAddNode::duplicate()
|
||||
{
|
||||
@@ -132,6 +205,12 @@ SHAMapAddNode::invalid()
|
||||
return SHAMapAddNode(0, 1, 0);
|
||||
}
|
||||
|
||||
inline void
|
||||
SHAMapAddNode::reset()
|
||||
{
|
||||
good_ = bad_ = duplicate_ = 0;
|
||||
}
|
||||
|
||||
inline std::string
|
||||
SHAMapAddNode::get() const
|
||||
{
|
||||
@@ -160,4 +239,14 @@ SHAMapAddNode::get() const
|
||||
return ret;
|
||||
}
|
||||
|
||||
inline SHAMapAddNode&
|
||||
SHAMapAddNode::operator+=(SHAMapAddNode const& n)
|
||||
{
|
||||
good_ += n.good_;
|
||||
bad_ += n.bad_;
|
||||
duplicate_ += n.duplicate_;
|
||||
|
||||
return *this;
|
||||
}
|
||||
|
||||
} // namespace xrpl
|
||||
|
||||
@@ -31,7 +31,19 @@ private:
|
||||
*/
|
||||
TaggedPointer hashesAndChildren_;
|
||||
|
||||
std::uint32_t fullBelowGen_ = 0;
|
||||
// Pin that wrapping fullBelowGen_ in an atomic keeps the packed layout, and that the load
|
||||
// isFullBelow() does once per node of every walk stays lock-free.
|
||||
static_assert(std::atomic<std::uint32_t>::is_always_lock_free);
|
||||
static_assert(sizeof(std::atomic<std::uint32_t>) == sizeof(std::uint32_t));
|
||||
static_assert(alignof(std::atomic<std::uint32_t>) == alignof(std::uint32_t));
|
||||
|
||||
/**
|
||||
* Written from more than one thread, since canonicalization shares a node
|
||||
* between maps and a walk can run with the acquisition lock released (see
|
||||
* SHAMap::state_). Relaxed both ways, since a generation is only compared
|
||||
* for equality and the children it vouches for are published through lock_.
|
||||
*/
|
||||
std::atomic<std::uint32_t> fullBelowGen_ = 0;
|
||||
std::uint16_t isBranch_ = 0;
|
||||
|
||||
/**
|
||||
@@ -204,13 +216,13 @@ SHAMapInnerNode::getBranchCount() const
|
||||
inline bool
|
||||
SHAMapInnerNode::isFullBelow(std::uint32_t generation) const
|
||||
{
|
||||
return fullBelowGen_ == generation;
|
||||
return fullBelowGen_.load(std::memory_order_relaxed) == generation;
|
||||
}
|
||||
|
||||
inline void
|
||||
SHAMapInnerNode::setFullBelowGen(std::uint32_t gen)
|
||||
{
|
||||
fullBelowGen_ = gen;
|
||||
fullBelowGen_.store(gen, std::memory_order_relaxed);
|
||||
}
|
||||
|
||||
} // namespace xrpl
|
||||
|
||||
@@ -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;
|
||||
@@ -279,7 +297,8 @@ Ledger::Ledger(Ledger const& prevLedger, NetClock::time_point closeTime)
|
||||
}
|
||||
|
||||
Ledger::Ledger(LedgerHeader const& info, Rules rules, Family& family)
|
||||
: immutable_(true)
|
||||
: immutable_(false)
|
||||
, mapHashesFromHeader_(true)
|
||||
, txMap_(SHAMapType::TRANSACTION, info.txHash, family)
|
||||
, stateMap_(SHAMapType::STATE, info.accountHash, family)
|
||||
, rules_(std::move(rules))
|
||||
@@ -308,27 +327,48 @@ 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)
|
||||
// An immutable ledger is treated as persistable, so an invalid map must never become
|
||||
// immutable. The validity test runs before anything is written, so a refusal leaves
|
||||
// the header as it was.
|
||||
if (!mapsValid())
|
||||
return false;
|
||||
|
||||
// Read now but written to the header below: getHash() can unshare a dirty tree, so it must run
|
||||
// while the map is still mutable, while the header may only be written once both maps are
|
||||
// settled. Skipped when the ledger is already immutable or when the header supplied these.
|
||||
bool const deriveMapHashes = !immutable_ && !mapHashesFromHeader_ && rehash;
|
||||
UInt256 const txHash = deriveMapHashes ? txMap_.getHash().asUInt256() : UInt256{};
|
||||
UInt256 const accountHash = deriveMapHashes ? stateMap_.getHash().asUInt256() : UInt256{};
|
||||
|
||||
// A concurrent walk can invalidate a map between the check above and here (see
|
||||
// SHAMap::state_), so the result is checked. That narrows the window rather than closing it,
|
||||
// since setInvalid() outranks Immutable.
|
||||
bool const bothImmutable = setMapsImmutable();
|
||||
SOMETIMES(!bothImmutable, "xrpl::Ledger::setImmutable : map invalidated while going immutable");
|
||||
if (!bothImmutable)
|
||||
return false; // LCOV_EXCL_LINE: only the walk named above reaches this, so no test does
|
||||
|
||||
// Written only now, so failing the check above leaves the header describing what the ledger was
|
||||
// built from. Forced, 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 +380,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
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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)
|
||||
@@ -143,12 +162,12 @@ SHAMap::visitDifferences(
|
||||
if (!function(*node))
|
||||
return;
|
||||
|
||||
// Nibbles run out at kLeafDepth, so only a leaf belongs there. A well-formed map never
|
||||
// holds an inner node at that depth: addKnownNode marks the map invalid rather than hooking
|
||||
// one in, and fetch-pack data is hash-verified against a validated root, so reaching this
|
||||
// means a defect or a corrupt store, not something a peer can provoke. Report the node
|
||||
// anyway - the wire form carries no depth, and the recipient hooks blobs in by hash - but
|
||||
// skip the children rather than letting getChildNodeID throw on them.
|
||||
// Nibbles run out at kLeafDepth, so only a leaf belongs there. addKnownNode marks the map
|
||||
// invalid on meeting an inner node at that depth, and fetch-pack data is hash-verified
|
||||
// against a validated root, so reaching this means a defect or a corrupt store. The
|
||||
// node is still reported, since the wire form carries no depth and the recipient hooks
|
||||
// blobs in by hash. Its children are skipped, since getChildNodeID has no answer past
|
||||
// kLeafDepth.
|
||||
if (nodeID.getDepth() >= kLeafDepth)
|
||||
{
|
||||
// LCOV_EXCL_START
|
||||
@@ -207,7 +226,12 @@ SHAMap::gmnProcessNodes(MissingNodes& mn, MissingNodes::StackEntry& se)
|
||||
// we already know this child node is missing
|
||||
fullBelow = false;
|
||||
}
|
||||
else if (!backed_ || !f_.getFullBelowCache()->touchIfExists(childHash.asUInt256()))
|
||||
// The depth test precedes the cache lookup for the same reason it does in addKnownNode():
|
||||
// the cache is keyed by node hash and shared across maps, and a hash covers a node's
|
||||
// children but not its depth. Skipping the shortcut forgoes an optimization only.
|
||||
else if (
|
||||
!backed_ || isLeafDepth(nodeID.getDepth() + 1) ||
|
||||
!f_.getFullBelowCache()->touchIfExists(childHash.asUInt256()))
|
||||
{
|
||||
bool pending = false;
|
||||
auto d = descendAsync(
|
||||
@@ -238,6 +262,17 @@ SHAMap::gmnProcessNodes(MissingNodes& mn, MissingNodes::StackEntry& se)
|
||||
if (--mn.max <= 0)
|
||||
return;
|
||||
}
|
||||
else if (d->isInner() && isLeafDepth(nodeID.getDepth() + 1))
|
||||
{
|
||||
// Only a leaf belongs that deep (see isLeafDepth and SHAMap::addKnownNode). A node
|
||||
// resolved locally reaches the walk without passing through addKnownNode(), so the
|
||||
// walk reaches this verdict itself. Ordered ahead of the full-below test below,
|
||||
// which canonicalization shares across maps.
|
||||
JLOG(journal_.warn()) << "Inner node at branch " << branch << " below " << nodeID
|
||||
<< " makes the map invalid";
|
||||
setInvalid();
|
||||
return;
|
||||
}
|
||||
else if (d->isInner() && !safeDowncast<SHAMapInnerNode*>(d)->isFullBelow(mn.generation))
|
||||
{
|
||||
mn.stack.push(se);
|
||||
@@ -301,6 +336,8 @@ SHAMap::gmnProcessDeferredReads(MissingNodes& mn)
|
||||
}
|
||||
else if ((mn.max > 0) && (mn.missingHashes.insert(nodeHash).second))
|
||||
{
|
||||
// getChildNodeID() is safe here: gmnProcessNodes refuses to descend into an inner node
|
||||
// at kLeafDepth, so a deferred parent sits at most one level above it.
|
||||
mn.missingNodes.emplace_back(parentID.getChildNodeID(branch), nodeHash.asUInt256());
|
||||
--mn.max;
|
||||
}
|
||||
@@ -311,17 +348,20 @@ SHAMap::gmnProcessDeferredReads(MissingNodes& mn)
|
||||
mn.deferred = 0;
|
||||
}
|
||||
|
||||
/**
|
||||
* Get a list of node IDs and hashes for nodes that are part of this SHAMap
|
||||
* but not available locally. The filter can hold alternate sources of
|
||||
* nodes that are not permanently stored locally
|
||||
*/
|
||||
std::vector<std::pair<SHAMapNodeID, UInt256>>
|
||||
SHAMap::getMissingNodes(int max, SHAMapSyncFilter const* filter)
|
||||
{
|
||||
XRPL_ASSERT(root_->getHash().isNonZero(), "xrpl::SHAMap::getMissingNodes : nonzero root hash");
|
||||
XRPL_ASSERT(max > 0, "xrpl::SHAMap::getMissingNodes : valid max input");
|
||||
|
||||
if (!isValid())
|
||||
{
|
||||
// The root node's own hash, since getHash() unshares the tree on a zero hash.
|
||||
JLOG(journal_.warn()) << "getMissingNodes called on an invalid map, root hash "
|
||||
<< root_->getHash() << " seq " << ledgerSeq();
|
||||
return {};
|
||||
}
|
||||
|
||||
MissingNodes mn(
|
||||
max,
|
||||
filter,
|
||||
@@ -354,6 +394,11 @@ SHAMap::getMissingNodes(int max, SHAMapSyncFilter const* filter)
|
||||
{
|
||||
gmnProcessNodes(mn, pos);
|
||||
|
||||
// The walk just invalidated the map. The loop stops descending here but falls through
|
||||
// to the drain below, since every posted read must be drained while `mn` is alive.
|
||||
if (!isValid())
|
||||
break;
|
||||
|
||||
if (mn.max <= 0)
|
||||
break;
|
||||
|
||||
@@ -383,6 +428,11 @@ SHAMap::getMissingNodes(int max, SHAMapSyncFilter const* filter)
|
||||
if (mn.deferred != 0)
|
||||
gmnProcessDeferredReads(mn);
|
||||
|
||||
// Reads are drained, so the map can be abandoned. What was collected belongs to a tree
|
||||
// that cannot exist.
|
||||
if (!isValid())
|
||||
return {};
|
||||
|
||||
if (mn.max <= 0)
|
||||
return std::move(mn.missingNodes);
|
||||
|
||||
@@ -416,6 +466,11 @@ SHAMap::getMissingNodes(int max, SHAMapSyncFilter const* filter)
|
||||
|
||||
} while (node != nullptr);
|
||||
|
||||
// addKnownNode() on another thread can write the verdict after the loop's own test, so the map
|
||||
// is judged once more before the result is returned.
|
||||
if (!isValid())
|
||||
return {}; // LCOV_EXCL_LINE: only that other thread reaches this, so no test does
|
||||
|
||||
if (mn.missingNodes.empty())
|
||||
clearSynching();
|
||||
|
||||
@@ -527,11 +582,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 +618,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 +661,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 +685,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. The node is reported as bad data.
|
||||
//
|
||||
// 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.
|
||||
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 +723,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 +842,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");
|
||||
@@ -859,10 +935,9 @@ SHAMap::verifyProofPath(UInt256 const& rootHash, UInt256 const& key, std::vector
|
||||
}
|
||||
else
|
||||
{
|
||||
// The hash chain up to rootHash only proves this leaf sits where the path claims,
|
||||
// not that it is the leaf for `key`: a peer could substitute any other leaf whose
|
||||
// subtree hashes to the same value at every level above it. Checking the terminal
|
||||
// leaf's own key is what ties the proof to `key` specifically.
|
||||
// The hash chain up to rootHash proves this leaf sits where the path claims. Any
|
||||
// leaf whose subtree hashes the same at every level above satisfies that chain, so
|
||||
// the terminal leaf's own key is what ties the proof to `key`.
|
||||
if (leafKey(*node) != key)
|
||||
return false;
|
||||
|
||||
|
||||
365
src/test/app/AcquireTestHelpers.h
Normal file
365
src/test/app/AcquireTestHelpers.h
Normal file
@@ -0,0 +1,365 @@
|
||||
#pragma once
|
||||
|
||||
#include <test/jtx/PeerStub.h>
|
||||
|
||||
#include <xrpld/app/ledger/ConsensusTransSetSF.h>
|
||||
#include <xrpld/overlay/Peer.h>
|
||||
#include <xrpld/overlay/PeerSet.h>
|
||||
|
||||
#include <xrpl/basics/SHAMapHash.h>
|
||||
#include <xrpl/basics/base_uint.h>
|
||||
#include <xrpl/protocol/Serializer.h>
|
||||
#include <xrpl/resource/Charge.h>
|
||||
#include <xrpl/shamap/SHAMapAddNode.h>
|
||||
#include <xrpl/shamap/SHAMapNodeID.h>
|
||||
#include <xrpl/shamap/SHAMapTreeNode.h>
|
||||
|
||||
#include <tests/libxrpl/shamap/DeepChain.h>
|
||||
|
||||
#include <xrpl.pb.h>
|
||||
|
||||
#include <atomic>
|
||||
#include <chrono>
|
||||
#include <cstddef>
|
||||
#include <cstdint>
|
||||
#include <functional>
|
||||
#include <memory>
|
||||
#include <mutex>
|
||||
#include <optional>
|
||||
#include <set>
|
||||
#include <string>
|
||||
#include <thread>
|
||||
#include <utility>
|
||||
#include <vector>
|
||||
|
||||
namespace xrpl::test {
|
||||
|
||||
// The chain builder needs only libxrpl, so it is shared with the gtest suites; see the header for
|
||||
// what keeps the protobuf reply builder below in this tree.
|
||||
using tests::DeepChain;
|
||||
|
||||
// A smallest-possible leaf plus its 4-byte HashPrefix lands one byte short of the floor
|
||||
// ConsensusTransSetSF::gotNode() parses at, so a chain's leaf stays below the parse threshold.
|
||||
static_assert(
|
||||
sizeof(std::uint32_t) + DeepChain::kLeafItemBytes < ConsensusTransSetSF::kMinTxNodeBytesToParse,
|
||||
"a smallest-possible leaf must stay below the resubmission floor");
|
||||
|
||||
/**
|
||||
* A peer that records what it was charged. Every other method comes from
|
||||
* PeerStub.
|
||||
*
|
||||
* One instance per packet keeps charges() unambiguous about which packet was
|
||||
* charged what.
|
||||
*
|
||||
* charges_ is unguarded: charging happens on the packet path, so every charge
|
||||
* lands on the thread that fed the packet in.
|
||||
*/
|
||||
class ChargeRecordingPeer : public PeerStub
|
||||
{
|
||||
public:
|
||||
/**
|
||||
* @param hasTxSet What hasTxSet() reports, which is how an acquisition
|
||||
* decides whether this peer is worth asking. Defaults to true, so a
|
||||
* peer handed straight to takeNodes() needs no argument.
|
||||
*/
|
||||
explicit ChargeRecordingPeer(bool hasTxSet = true) : PeerStub(nextId()), hasTxSet_(hasTxSet)
|
||||
{
|
||||
}
|
||||
|
||||
void
|
||||
charge(resource::Charge const& fee, std::string const&) override
|
||||
{
|
||||
charges_.push_back(fee);
|
||||
}
|
||||
|
||||
/**
|
||||
* @return Every fee this peer has been charged, in the order charged.
|
||||
*/
|
||||
[[nodiscard]] std::vector<resource::Charge> const&
|
||||
charges() const
|
||||
{
|
||||
return charges_;
|
||||
}
|
||||
|
||||
// PeerStub returns false for both, and an acquisition asks the peers reporting true.
|
||||
|
||||
[[nodiscard]] bool
|
||||
hasTxSet(UInt256 const&) const override
|
||||
{
|
||||
return hasTxSet_;
|
||||
}
|
||||
|
||||
[[nodiscard]] bool
|
||||
hasLedger(UInt256 const&, std::uint32_t) const override
|
||||
{
|
||||
return true;
|
||||
}
|
||||
|
||||
private:
|
||||
/**
|
||||
* The next id to hand out, distinct per instance because
|
||||
* RequestCountingPeerSet dedups by tracked id and PeerStub's own default is
|
||||
* zero for every instance.
|
||||
*
|
||||
* @return The id.
|
||||
*/
|
||||
[[nodiscard]] static ID
|
||||
nextId()
|
||||
{
|
||||
static std::atomic<ID> next{1};
|
||||
return next++;
|
||||
}
|
||||
|
||||
std::vector<resource::Charge> charges_;
|
||||
bool hasTxSet_;
|
||||
};
|
||||
|
||||
/**
|
||||
* A peer set that counts the requests an acquisition makes through it. Offers
|
||||
* peers to a hasItem/onPeerAdded callback pair, hard-filtered by hasItem (which
|
||||
* only scores in the real peer set) and deduped by tracked id.
|
||||
*
|
||||
* A count is a call, not a delivery: a request naming no peer reaches nobody
|
||||
* while this set tracks none. requests() and broadcasts() are counted apart.
|
||||
*
|
||||
* Every write here is guarded, since the retry timer drives addPeers() and
|
||||
* sendRequest() from a job thread while the test reads the results.
|
||||
*/
|
||||
class RequestCountingPeerSet : public PeerSet
|
||||
{
|
||||
public:
|
||||
/**
|
||||
* @param candidates The peers addPeers() may offer, in the order they are
|
||||
* considered. Fixed at construction. Empty offers no one.
|
||||
*/
|
||||
explicit RequestCountingPeerSet(std::vector<std::shared_ptr<Peer>> candidates = {})
|
||||
: candidates_(std::move(candidates))
|
||||
{
|
||||
}
|
||||
|
||||
/**
|
||||
* Offer the candidates to the caller, the way the real peer set offers the
|
||||
* peers the overlay is tracking.
|
||||
*
|
||||
* @param limit The most peers to add, recorded for firstLimit().
|
||||
* @param hasItem Hard-filters the candidates worth asking, where the real
|
||||
* peer set only scores with it.
|
||||
* @param onPeerAdded Called for each selected candidate.
|
||||
*/
|
||||
void
|
||||
addPeers(
|
||||
std::size_t limit,
|
||||
std::function<bool(std::shared_ptr<Peer> const&)> hasItem,
|
||||
std::function<void(std::shared_ptr<Peer> const&)> onPeerAdded) override
|
||||
{
|
||||
std::vector<std::shared_ptr<Peer>> selected;
|
||||
{
|
||||
std::scoped_lock const lock(mutex_);
|
||||
|
||||
if (!firstLimit_)
|
||||
firstLimit_ = limit;
|
||||
|
||||
for (auto const& candidate : candidates_)
|
||||
{
|
||||
if (selected.size() >= limit)
|
||||
break;
|
||||
// Dedup by tracked id, like the real peer set: onPeerAdded runs once per candidate.
|
||||
if (hasItem(candidate) && addedPeers_.insert(candidate->id()).second)
|
||||
selected.push_back(candidate);
|
||||
}
|
||||
}
|
||||
|
||||
// Outside the lock: onPeerAdded() reenters this object through sendRequest().
|
||||
for (auto const& peer : selected)
|
||||
onPeerAdded(peer);
|
||||
}
|
||||
|
||||
/**
|
||||
* Record the request, and record it in the broadcast count as well when it
|
||||
* names no peer.
|
||||
*
|
||||
* @param peer The peer to ask, or null to ask every tracked peer.
|
||||
*/
|
||||
void
|
||||
sendRequest(
|
||||
::google::protobuf::Message const&,
|
||||
protocol::MessageType,
|
||||
std::shared_ptr<Peer> const& peer) override
|
||||
{
|
||||
std::scoped_lock const lock(mutex_);
|
||||
++requests_;
|
||||
if (!peer)
|
||||
++broadcasts_;
|
||||
}
|
||||
|
||||
/**
|
||||
* The ids of every peer addPeers() has selected, which is what an
|
||||
* acquisition takes for the peers it is tracking.
|
||||
*
|
||||
* Unguarded: every caller of this and of addPeers() is an acquisition
|
||||
* holding its own mtx_, so the ids stay fixed while a caller iterates. A
|
||||
* test thread reads addedPeers() instead.
|
||||
* InboundLedger::getPeerCount() reports zero for these, since it resolves
|
||||
* ids through the overlay, while these peers live only in this harness.
|
||||
*
|
||||
* @return The ids.
|
||||
*/
|
||||
[[nodiscard]] std::set<Peer::ID> const&
|
||||
getPeerIds() const override
|
||||
{
|
||||
return addedPeers_;
|
||||
}
|
||||
|
||||
/**
|
||||
* @return How many requests have been sent through this peer set, counting
|
||||
* a broadcast as one.
|
||||
*/
|
||||
[[nodiscard]] int
|
||||
requests() const
|
||||
{
|
||||
std::scoped_lock const lock(mutex_);
|
||||
return requests_;
|
||||
}
|
||||
|
||||
/**
|
||||
* How many of those requests carried no peer of their own.
|
||||
*
|
||||
* @return The count.
|
||||
*/
|
||||
[[nodiscard]] int
|
||||
broadcasts() const
|
||||
{
|
||||
std::scoped_lock const lock(mutex_);
|
||||
return broadcasts_;
|
||||
}
|
||||
|
||||
/**
|
||||
* The limit the first addPeers() call asked for, which is init()'s, since
|
||||
* onTimer() keeps calling addPeers(1) for as long as an acquisition runs.
|
||||
*
|
||||
* @return The limit, or nullopt if addPeers() has not been called.
|
||||
*/
|
||||
[[nodiscard]] std::optional<std::size_t>
|
||||
firstLimit() const
|
||||
{
|
||||
std::scoped_lock const lock(mutex_);
|
||||
return firstLimit_;
|
||||
}
|
||||
|
||||
/**
|
||||
* A set, because onTimer() keeps re-offering the same candidates.
|
||||
*
|
||||
* @return The ids of every peer addPeers() has selected so far.
|
||||
*/
|
||||
[[nodiscard]] std::set<Peer::ID>
|
||||
addedPeers() const
|
||||
{
|
||||
std::scoped_lock const lock(mutex_);
|
||||
return addedPeers_;
|
||||
}
|
||||
|
||||
private:
|
||||
std::vector<std::shared_ptr<Peer>> const candidates_;
|
||||
|
||||
mutable std::mutex mutex_;
|
||||
int requests_{0};
|
||||
int broadcasts_{0};
|
||||
std::optional<std::size_t> firstLimit_;
|
||||
std::set<Peer::ID> addedPeers_;
|
||||
};
|
||||
|
||||
/**
|
||||
* The given nodes of a chain as a TMLedgerData, so a test can go through the
|
||||
* real dispatch. Not in DeepChain, since the protobuf types are xrpld and that
|
||||
* header is shared with the libxrpl-only gtest binary.
|
||||
*
|
||||
* @param chain The chain the nodes came from, which names the reply by default.
|
||||
* @param data The nodes to include, each with its claimed position.
|
||||
* @param type The reply type, which selects which map the receiver applies it
|
||||
* to.
|
||||
* @param ledgerHash The hash the reply claims to be about, defaulting to the
|
||||
* chain root for a TX set. A ledger acquisition wants its header hash
|
||||
* here, since the chain root is only that ledger's account hash.
|
||||
* @param ledgerSeq The sequence to name in the reply.
|
||||
* @return The reply packet.
|
||||
*/
|
||||
[[nodiscard]] inline std::shared_ptr<protocol::TMLedgerData>
|
||||
packetFor(
|
||||
DeepChain const& chain,
|
||||
std::vector<std::pair<SHAMapNodeID, SHAMapTreeNodePtr>> const& data,
|
||||
protocol::TMLedgerInfoType type = protocol::liTS_CANDIDATE,
|
||||
std::optional<UInt256> const& ledgerHash = std::nullopt,
|
||||
std::uint32_t ledgerSeq = 0)
|
||||
{
|
||||
auto packet = std::make_shared<protocol::TMLedgerData>();
|
||||
auto const hash = ledgerHash.value_or(chain.rootHash.asUInt256());
|
||||
packet->set_ledgerhash(hash.data(), UInt256::size());
|
||||
packet->set_ledgerseq(ledgerSeq);
|
||||
packet->set_type(type);
|
||||
|
||||
for (auto const& [nodeID, node] : data)
|
||||
{
|
||||
Serializer s;
|
||||
node->serializeForWire(s);
|
||||
|
||||
auto* const ledgerNode = packet->add_nodes();
|
||||
ledgerNode->set_nodedata(s.peekData().data(), s.peekData().size());
|
||||
|
||||
// A leaf carries its own key, so the receiver rebuilds its position from that plus a
|
||||
// depth. An inner node has no key and needs the full ID. The two fields are a oneof.
|
||||
if (node->isLeaf())
|
||||
{
|
||||
ledgerNode->set_depth(nodeID.getDepth());
|
||||
}
|
||||
else
|
||||
{
|
||||
ledgerNode->set_id(nodeID.getRawString());
|
||||
}
|
||||
}
|
||||
|
||||
return packet;
|
||||
}
|
||||
|
||||
/**
|
||||
* Poll until the condition holds, or give up. An acquisition's own timer and
|
||||
* the jobs it hands finished work to both run on other threads.
|
||||
*
|
||||
* @param condition What to wait for.
|
||||
* @param deadline The longest to wait.
|
||||
* @return Whether the condition held before the deadline.
|
||||
*/
|
||||
[[nodiscard]] inline bool
|
||||
waitFor(
|
||||
std::function<bool()> const& condition,
|
||||
std::chrono::steady_clock::duration deadline = std::chrono::seconds{10})
|
||||
{
|
||||
auto const giveUp = std::chrono::steady_clock::now() + deadline;
|
||||
while (std::chrono::steady_clock::now() < giveUp)
|
||||
{
|
||||
if (condition())
|
||||
return true;
|
||||
std::this_thread::sleep_for(std::chrono::milliseconds{10});
|
||||
}
|
||||
return condition();
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether a batch verdict carries exactly the given counts.
|
||||
*
|
||||
* The counts, since get() is a log format. It is pinned once, in the
|
||||
* SHAMapAddNode tests, and is what to pass BEAST_EXPECTS() as the reason a
|
||||
* check here failed.
|
||||
*
|
||||
* @param san The verdict to check.
|
||||
* @param good How many nodes the batch should have hooked in.
|
||||
* @param bad How many it should have rejected.
|
||||
* @param duplicate How many it should have already held.
|
||||
* @return Whether the verdict matches.
|
||||
*/
|
||||
[[nodiscard]] inline bool
|
||||
tallyIs(SHAMapAddNode const& san, int good, int bad, int duplicate)
|
||||
{
|
||||
return san.getGood() == good && san.getBad() == bad && san.getDuplicate() == duplicate;
|
||||
}
|
||||
|
||||
} // namespace xrpl::test
|
||||
996
src/test/app/InboundLedger_test.cpp
Normal file
996
src/test/app/InboundLedger_test.cpp
Normal file
@@ -0,0 +1,996 @@
|
||||
#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/Blob.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);
|
||||
}
|
||||
|
||||
/**
|
||||
* The same, as the timer chain does.
|
||||
*/
|
||||
void
|
||||
triggerTimeout()
|
||||
{
|
||||
trigger(nullptr, TriggerReason::Timeout);
|
||||
}
|
||||
|
||||
/**
|
||||
* Record how many timeouts have elapsed.
|
||||
*
|
||||
* @param timeouts The count to record.
|
||||
*/
|
||||
void
|
||||
setTimeouts(int timeouts)
|
||||
{
|
||||
ScopedLockType const sl(mtx_);
|
||||
timeouts_ = timeouts;
|
||||
}
|
||||
|
||||
/**
|
||||
* Forget any recorded progress.
|
||||
*/
|
||||
void
|
||||
clearProgress()
|
||||
{
|
||||
ScopedLockType const sl(mtx_);
|
||||
progress_ = false;
|
||||
}
|
||||
|
||||
/**
|
||||
* 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 completed by a walk rather than by tryDB() is settled before it
|
||||
* is reported complete.
|
||||
*
|
||||
* The order is visible only to a reader of isComplete(), which is read
|
||||
* without mtx_, so this case pins the invariant instead: every ledger the
|
||||
* acquisition reports is immutable.
|
||||
*
|
||||
* @param env The environment to run in.
|
||||
*/
|
||||
void
|
||||
testWalkSettlesBeforeReportingComplete(jtx::Env& env)
|
||||
{
|
||||
testcase("A ledger completed by a walk is settled before it is reported");
|
||||
|
||||
auto const chain = DeepChain::toLeaf(2, nextSeed());
|
||||
auto const header = makeHeader(chain);
|
||||
|
||||
// Everything except the leaf, so tryDB() can root the state map but its walk still has
|
||||
// something to ask for, which leaves the completion to trigger().
|
||||
storeHeader(env, header);
|
||||
storeStateNodes(env, header, chain, chain.deepestDepth - 1);
|
||||
|
||||
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->isComplete());
|
||||
BEAST_EXPECT(!acquire->isFailed());
|
||||
|
||||
// The walk can finish only now, which places the completion in the walk.
|
||||
storeStateNodes(env, header, chain, chain.deepestDepth);
|
||||
|
||||
acquire->triggerAdded();
|
||||
|
||||
BEAST_EXPECT(acquire->isComplete());
|
||||
BEAST_EXPECT(!acquire->isFailed());
|
||||
|
||||
auto const settled = acquire->getLedger();
|
||||
BEAST_EXPECT(settled != nullptr);
|
||||
if (settled)
|
||||
BEAST_EXPECT(settled->isImmutable());
|
||||
|
||||
// done()'s success arm ran, so the ledger reached LedgerMaster and the hash is absent
|
||||
// from the failure list.
|
||||
BEAST_EXPECT(env.app().getLedgerMaster().getLedgerByHash(header.hash) != nullptr);
|
||||
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 ledger assembled from local data must be judged even when only
|
||||
* one map is settled.
|
||||
*
|
||||
* tryDB() walks both maps to see what is on hand, and a fetch pack is
|
||||
* checked against each node's own hash rather than the shape it
|
||||
* implies, so a whole chain can resolve locally without passing
|
||||
* through addKnownNode().
|
||||
*
|
||||
* The asymmetry is the point: the transaction map is the chain, so
|
||||
* its walk abandons it, while the state root is a hash no fetch pack
|
||||
* supplies, leaving that map merely incomplete. tryDB() therefore
|
||||
* sets neither flag and has to reach the verdict itself, since the
|
||||
* setImmutable() call further down needs both.
|
||||
*
|
||||
* @param env The environment to run in.
|
||||
*/
|
||||
void
|
||||
testLocalChainFailsAcquire(jtx::Env& env)
|
||||
{
|
||||
testcase("A chain found locally fails the acquire");
|
||||
|
||||
DeepChain const chain{nextSeed()};
|
||||
|
||||
// The chain as the transaction root, and an arbitrary hash, seeded nowhere, as the state
|
||||
// root.
|
||||
auto const header = makeHeader(chain.rootHash.asUInt256(), UInt256{99});
|
||||
auto& ledgerMaster = env.app().getLedgerMaster();
|
||||
|
||||
// The header, prefixed the way tryDB() expects to find it in a fetch pack.
|
||||
Serializer hs;
|
||||
hs.add32(HashPrefix::LedgerMaster);
|
||||
addRaw(header, hs);
|
||||
ledgerMaster.addFetchPack(header.hash, std::make_shared<Blob>(hs.modData()));
|
||||
|
||||
// Every node of the chain, keyed by its own hash. TransactionStateSF::getNode() reads
|
||||
// these, so the transaction-map walk resolves the whole chain with no peer involved.
|
||||
for (auto depth = 0u; depth <= SHAMap::kLeafDepth; ++depth)
|
||||
{
|
||||
ledgerMaster.addFetchPack(
|
||||
chain.nodeAt(depth)->getHash().asUInt256(),
|
||||
std::make_shared<Blob>(chain.prefixedNodeAt(depth)));
|
||||
}
|
||||
|
||||
auto acquire = std::make_shared<InboundLedger>(
|
||||
env.app(),
|
||||
header.hash,
|
||||
header.seq,
|
||||
InboundLedger::Reason::GENERIC,
|
||||
stopwatch(),
|
||||
std::make_unique<RequestCountingPeerSet>());
|
||||
|
||||
// checkLocal() routes into tryDB() with no peer data. It reports true only because the
|
||||
// acquisition ended.
|
||||
BEAST_EXPECT(acquire->checkLocal());
|
||||
|
||||
BEAST_EXPECT(acquire->isFailed());
|
||||
BEAST_EXPECT(!acquire->isComplete());
|
||||
|
||||
// A failed acquisition reports no ledger, though it still holds the partial one it built.
|
||||
BEAST_EXPECT(acquire->getLedger() == nullptr);
|
||||
}
|
||||
|
||||
/**
|
||||
* The aggressive-retry branch of trigger() must judge a map the walk
|
||||
* abandoned.
|
||||
*
|
||||
* That branch reads an empty getNeededHashes() result as a complete map,
|
||||
* and the walk it runs can reach the invalid verdict itself once nodes
|
||||
* resolve from local storage. Only the root is local when tryDB() runs, so
|
||||
* the state map holds a root and the walk stops one level down. The branch
|
||||
* also needs a timeout count above kLedgerBecomeAggressiveThreshold, which
|
||||
* the case records directly rather than waiting out the timer chain.
|
||||
*
|
||||
* @param env The environment to run in.
|
||||
*/
|
||||
void
|
||||
testAggressiveRetryJudgesLocalMap(jtx::Env& env)
|
||||
{
|
||||
testcase("An aggressive retry judges a map the walk abandoned");
|
||||
|
||||
DeepChain const chain{nextSeed()};
|
||||
|
||||
// The chain as the state root, and no transactions, so only the state map is in play.
|
||||
auto const header = makeHeader(chain);
|
||||
auto& ledgerMaster = env.app().getLedgerMaster();
|
||||
|
||||
Serializer hs;
|
||||
hs.add32(HashPrefix::LedgerMaster);
|
||||
addRaw(header, hs);
|
||||
ledgerMaster.addFetchPack(header.hash, std::make_shared<Blob>(hs.modData()));
|
||||
|
||||
// Only the root, so the state map gets a root but the walk stops one level down.
|
||||
ledgerMaster.addFetchPack(
|
||||
chain.nodeAt(0)->getHash().asUInt256(),
|
||||
std::make_shared<Blob>(chain.prefixedNodeAt(0)));
|
||||
|
||||
auto acquire = std::make_shared<TestableInboundLedger>(
|
||||
env.app(),
|
||||
header.hash,
|
||||
header.seq,
|
||||
InboundLedger::Reason::GENERIC,
|
||||
stopwatch(),
|
||||
std::make_unique<RequestCountingPeerSet>());
|
||||
|
||||
// The acquisition is alive: it has the header and a state root, and still wants the rest.
|
||||
BEAST_EXPECT(!acquire->checkLocal());
|
||||
BEAST_EXPECT(!acquire->isFailed());
|
||||
BEAST_EXPECT(acquire->getJson(0)[jss::have_header].asBool());
|
||||
BEAST_EXPECT(!acquire->getJson(0)[jss::have_state].asBool());
|
||||
|
||||
auto const ledger = mutableLedger(*acquire);
|
||||
BEAST_EXPECT(ledger != nullptr);
|
||||
if (!ledger)
|
||||
return;
|
||||
BEAST_EXPECT(ledger->stateMap().isValid());
|
||||
|
||||
// The rest of the chain becomes resolvable only now, which places the verdict in this walk.
|
||||
for (auto depth = 1u; depth <= SHAMap::kLeafDepth; ++depth)
|
||||
{
|
||||
ledgerMaster.addFetchPack(
|
||||
chain.nodeAt(depth)->getHash().asUInt256(),
|
||||
std::make_shared<Blob>(chain.prefixedNodeAt(depth)));
|
||||
}
|
||||
|
||||
// kLedgerBecomeAggressiveThreshold is 4 and file-local, so name the requirement here.
|
||||
acquire->setTimeouts(5);
|
||||
acquire->clearProgress();
|
||||
acquire->triggerTimeout();
|
||||
|
||||
// The walk resolved the chain locally and abandoned the map, and trigger() recorded that
|
||||
// rather than reading the empty result as a finished acquisition.
|
||||
BEAST_EXPECT(!ledger->stateMap().isValid());
|
||||
BEAST_EXPECT(acquire->isFailed());
|
||||
BEAST_EXPECT(!acquire->isComplete());
|
||||
|
||||
// have_state is the discriminating assertion: only this guard leaves it false, since every
|
||||
// have-flag is set on the way to the setImmutable() backstop in done().
|
||||
BEAST_EXPECT(!acquire->getJson(0)[jss::have_state].asBool());
|
||||
|
||||
// The same branch with no header, which is the other arm of hasInvalidMap(): no map, so
|
||||
// the arm answers false. getNeededHashes() has asked for the header, so the non-empty
|
||||
// branch is the right one and the acquisition stays alive.
|
||||
auto headerless = std::make_shared<TestableInboundLedger>(
|
||||
env.app(),
|
||||
UInt256{7},
|
||||
0,
|
||||
InboundLedger::Reason::GENERIC,
|
||||
stopwatch(),
|
||||
std::make_unique<RequestCountingPeerSet>());
|
||||
|
||||
headerless->setTimeouts(5);
|
||||
headerless->clearProgress();
|
||||
headerless->triggerTimeout();
|
||||
|
||||
BEAST_EXPECT(mutableLedger(*headerless) == nullptr);
|
||||
BEAST_EXPECT(!headerless->isFailed());
|
||||
BEAST_EXPECT(!headerless->isComplete());
|
||||
}
|
||||
|
||||
/**
|
||||
* The ordinary trigger() path must judge a map its own walk abandoned.
|
||||
*
|
||||
* Covers the state-map walk trigger() runs with mtx_ released, where the
|
||||
* case above covers the empty getNeededHashes() branch. That verdict has to
|
||||
* be read before the flags below it. Staged as that case is, with the
|
||||
* timeout count left at zero to stay off the aggressive branch.
|
||||
*
|
||||
* @param env The environment to run in.
|
||||
*/
|
||||
void
|
||||
testWalkJudgesMapOnOrdinaryTrigger(jtx::Env& env)
|
||||
{
|
||||
testcase("An ordinary trigger judges a map its walk abandoned");
|
||||
|
||||
DeepChain const chain{nextSeed()};
|
||||
|
||||
auto const header = makeHeader(chain);
|
||||
auto& ledgerMaster = env.app().getLedgerMaster();
|
||||
|
||||
Serializer hs;
|
||||
hs.add32(HashPrefix::LedgerMaster);
|
||||
addRaw(header, hs);
|
||||
ledgerMaster.addFetchPack(header.hash, std::make_shared<Blob>(hs.modData()));
|
||||
|
||||
// Only the root, so the state map gets a root but the walk stops one level down.
|
||||
ledgerMaster.addFetchPack(
|
||||
chain.nodeAt(0)->getHash().asUInt256(),
|
||||
std::make_shared<Blob>(chain.prefixedNodeAt(0)));
|
||||
|
||||
auto peerSet = std::make_unique<RequestCountingPeerSet>();
|
||||
auto* const peerSetPtr = peerSet.get();
|
||||
|
||||
auto acquire = std::make_shared<TestableInboundLedger>(
|
||||
env.app(),
|
||||
header.hash,
|
||||
header.seq,
|
||||
InboundLedger::Reason::GENERIC,
|
||||
stopwatch(),
|
||||
std::move(peerSet));
|
||||
|
||||
BEAST_EXPECT(!acquire->checkLocal());
|
||||
BEAST_EXPECT(!acquire->isFailed());
|
||||
BEAST_EXPECT(acquire->getJson(0)[jss::have_header].asBool());
|
||||
BEAST_EXPECT(!acquire->getJson(0)[jss::have_state].asBool());
|
||||
|
||||
auto const ledger = mutableLedger(*acquire);
|
||||
BEAST_EXPECT(ledger != nullptr);
|
||||
if (!ledger)
|
||||
return;
|
||||
BEAST_EXPECT(ledger->stateMap().isValid());
|
||||
|
||||
// The rest of the chain becomes resolvable only now, which places the verdict in this walk.
|
||||
for (auto depth = 1u; depth <= SHAMap::kLeafDepth; ++depth)
|
||||
{
|
||||
ledgerMaster.addFetchPack(
|
||||
chain.nodeAt(depth)->getHash().asUInt256(),
|
||||
std::make_shared<Blob>(chain.prefixedNodeAt(depth)));
|
||||
}
|
||||
|
||||
// Sound going in, which is the first half of the tripwire below.
|
||||
BEAST_EXPECT(ledger->stateMap().isValid());
|
||||
|
||||
int const requestsBefore = peerSetPtr->requests();
|
||||
acquire->triggerAdded();
|
||||
|
||||
// The walk resolved the chain locally and abandoned the map, and trigger() recorded that
|
||||
// rather than reading the empty node list as a finished state map.
|
||||
BEAST_EXPECT(!ledger->stateMap().isValid());
|
||||
BEAST_EXPECT(acquire->isFailed());
|
||||
|
||||
// have_state is the discriminating assertion. isComplete() is not: it reads false whether
|
||||
// this guard fired or the setImmutable() backstop in done() did.
|
||||
BEAST_EXPECT(!acquire->getJson(0)[jss::have_state].asBool());
|
||||
|
||||
// A tripwire that keeps this case on the branch it names: the aggressive branch needs both
|
||||
// a TIMEOUT reason and a count above kLedgerBecomeAggressiveThreshold.
|
||||
BEAST_EXPECT(acquire->getJson(0)[jss::timeouts].asInt() == 0);
|
||||
|
||||
// The request count is unchanged after the verdict, so the guard ended the round.
|
||||
BEAST_EXPECT(peerSetPtr->requests() == requestsBefore);
|
||||
}
|
||||
|
||||
/**
|
||||
* 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. Past
|
||||
* kLedgerTimeoutRetriesMax the acquisition fails itself and done() records
|
||||
* that, so the same ledger is not asked for again. The store is empty for
|
||||
* this hash and no reply arrives, so every tick counts a timeout.
|
||||
*
|
||||
* @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. This case reaches getLedger() with no ledger ever built, so
|
||||
// both arms of the failed_ gate answer null here. The local-chain case above is what
|
||||
// covers the gate itself: tryDB() builds a partial ledger and then fails, so the gate is
|
||||
// the only reason the answer is null.
|
||||
BEAST_EXPECT(acquire->getLedger() == nullptr);
|
||||
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);
|
||||
testWalkSettlesBeforeReportingComplete(env);
|
||||
testInvalidatedLedgerFailsInDone(env);
|
||||
testLocalFailureSignalsDone(env);
|
||||
testPeerZeroAccountHashFails(env);
|
||||
testPeerHeaderWithoutTransactionsCompletes(env);
|
||||
testLocalChainFailsAcquire(env);
|
||||
testAggressiveRetryJudgesLocalMap(env);
|
||||
testWalkJudgesMapOnOrdinaryTrigger(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
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
|
||||
1063
src/test/app/TransactionAcquire_test.cpp
Normal file
1063
src/test/app/TransactionAcquire_test.cpp
Normal file
File diff suppressed because it is too large
Load Diff
@@ -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));
|
||||
|
||||
@@ -1,123 +0,0 @@
|
||||
#pragma once
|
||||
|
||||
#include <xrpl/basics/ByteUtilities.h>
|
||||
#include <xrpl/basics/base_uint.h>
|
||||
#include <xrpl/basics/chrono.h>
|
||||
#include <xrpl/basics/contract.h>
|
||||
#include <xrpl/beast/utility/Journal.h>
|
||||
#include <xrpl/config/BasicConfig.h>
|
||||
#include <xrpl/config/Constants.h>
|
||||
#include <xrpl/nodestore/Database.h>
|
||||
#include <xrpl/nodestore/DummyScheduler.h>
|
||||
#include <xrpl/nodestore/Manager.h>
|
||||
#include <xrpl/shamap/Family.h>
|
||||
#include <xrpl/shamap/FullBelowCache.h>
|
||||
#include <xrpl/shamap/TreeNodeCache.h>
|
||||
|
||||
#include <cstdint>
|
||||
#include <memory>
|
||||
#include <stdexcept>
|
||||
|
||||
namespace xrpl::test {
|
||||
|
||||
/**
|
||||
* Test implementation of Family for unit tests.
|
||||
*
|
||||
* Uses an in-memory NodeStore database and simple caches.
|
||||
* The missingNode methods throw since tests shouldn't encounter missing nodes.
|
||||
*/
|
||||
class TestFamily : public Family
|
||||
{
|
||||
private:
|
||||
std::unique_ptr<node_store::Database> db_;
|
||||
TestStopwatch clock_;
|
||||
std::shared_ptr<FullBelowCache> fbCache_;
|
||||
std::shared_ptr<TreeNodeCache> tnCache_;
|
||||
node_store::DummyScheduler scheduler_;
|
||||
beast::Journal j_;
|
||||
|
||||
public:
|
||||
explicit TestFamily(beast::Journal j)
|
||||
: fbCache_(std::make_shared<FullBelowCache>("TestFamily full below cache", clock_, j))
|
||||
, tnCache_(
|
||||
std::make_shared<TreeNodeCache>(
|
||||
"TestFamily tree node cache",
|
||||
65536,
|
||||
std::chrono::minutes{1},
|
||||
clock_,
|
||||
j))
|
||||
, j_(j)
|
||||
{
|
||||
Section config;
|
||||
config.set(Keys::kType, "memory");
|
||||
config.set(Keys::kPath, "TestFamily");
|
||||
db_ = node_store::Manager::instance().makeDatabase(megabytes(4), scheduler_, 1, config, j);
|
||||
}
|
||||
|
||||
node_store::Database&
|
||||
db() override
|
||||
{
|
||||
return *db_;
|
||||
}
|
||||
|
||||
[[nodiscard]] node_store::Database const&
|
||||
db() const override
|
||||
{
|
||||
return *db_;
|
||||
}
|
||||
|
||||
beast::Journal const&
|
||||
journal() override
|
||||
{
|
||||
return j_;
|
||||
}
|
||||
|
||||
std::shared_ptr<FullBelowCache>
|
||||
getFullBelowCache() override
|
||||
{
|
||||
return fbCache_;
|
||||
}
|
||||
|
||||
std::shared_ptr<TreeNodeCache>
|
||||
getTreeNodeCache() override
|
||||
{
|
||||
return tnCache_;
|
||||
}
|
||||
|
||||
void
|
||||
sweep() override
|
||||
{
|
||||
fbCache_->sweep();
|
||||
tnCache_->sweep();
|
||||
}
|
||||
|
||||
void
|
||||
missingNodeAcquireBySeq(std::uint32_t refNum, UInt256 const& nodeHash) override
|
||||
{
|
||||
Throw<std::runtime_error>("TestFamily: missing node (by seq)");
|
||||
}
|
||||
|
||||
void
|
||||
missingNodeAcquireByHash(UInt256 const& refHash, std::uint32_t refNum) override
|
||||
{
|
||||
Throw<std::runtime_error>("TestFamily: missing node (by hash)");
|
||||
}
|
||||
|
||||
void
|
||||
reset() override
|
||||
{
|
||||
(*fbCache_).reset();
|
||||
(*tnCache_).reset();
|
||||
}
|
||||
|
||||
/**
|
||||
* Access the test clock for time manipulation in tests.
|
||||
*/
|
||||
TestStopwatch&
|
||||
clock()
|
||||
{
|
||||
return clock_;
|
||||
}
|
||||
};
|
||||
|
||||
} // namespace xrpl::test
|
||||
@@ -12,8 +12,8 @@
|
||||
|
||||
#include <boost/asio/io_context.hpp>
|
||||
|
||||
#include <helpers/TestFamily.h>
|
||||
#include <helpers/TestSink.h>
|
||||
#include <shamap/common.h>
|
||||
|
||||
#include <cstdint>
|
||||
#include <memory>
|
||||
@@ -73,7 +73,7 @@ class TestServiceRegistry : public ServiceRegistry
|
||||
{
|
||||
TestLogs logs_{beast::Severity::Warning};
|
||||
boost::asio::io_context ioContext_;
|
||||
TestFamily family_{logs_.journal("TestFamily")};
|
||||
tests::TestNodeFamily family_{logs_.journal("TestNodeFamily")};
|
||||
LoadFeeTrack feeTrack_{logs_.journal("LoadFeeTrack")};
|
||||
TestNetworkIDService networkIDService_;
|
||||
HashRouter hashRouter_{HashRouter::Setup{}, stopwatch()};
|
||||
|
||||
@@ -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;
|
||||
|
||||
|
||||
306
src/tests/libxrpl/shamap/DeepChain.h
Normal file
306
src/tests/libxrpl/shamap/DeepChain.h
Normal file
@@ -0,0 +1,306 @@
|
||||
#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 and, under withDecoys(), an unresolvable second one.
|
||||
*
|
||||
* 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 chain's nodes. Two chains built from one seed hold
|
||||
* the same nodes, which caches and fetch packs key by hash.
|
||||
*/
|
||||
explicit DeepChain(unsigned int seed = 1) : DeepChain(std::nullopt, seed, Decoy::No)
|
||||
{
|
||||
}
|
||||
|
||||
/**
|
||||
* The same chain, with an unresolvable second child at every level, so a
|
||||
* backed map's descendAsync() posts a real asynchronous read per level.
|
||||
*
|
||||
* Offered only for this shape: the decoy sits on branch 1, which is free
|
||||
* only while pathKey is zero.
|
||||
*
|
||||
* @param seed Varies the whole chain. See the constructor.
|
||||
* @return The chain.
|
||||
*/
|
||||
[[nodiscard]] static DeepChain
|
||||
withDecoys(unsigned int seed = 1)
|
||||
{
|
||||
return DeepChain{std::nullopt, seed, Decoy::Yes};
|
||||
}
|
||||
|
||||
/**
|
||||
* 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, Decoy::No};
|
||||
}
|
||||
|
||||
/**
|
||||
* 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:
|
||||
// Whether each level carries a second child that stays unresolvable.
|
||||
enum class Decoy { No, Yes };
|
||||
|
||||
/**
|
||||
* 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.
|
||||
* @param decoy Whether every level carries an unresolvable second child.
|
||||
*/
|
||||
DeepChain(std::optional<unsigned int> leafDepth, unsigned int seed, Decoy decoy)
|
||||
: 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}}, decoy);
|
||||
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(), decoy);
|
||||
}
|
||||
|
||||
/**
|
||||
* Fill in inner nodes, each with one real child (and, under Decoy::Yes, an
|
||||
* additional unresolvable second 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.
|
||||
* @param decoy Whether to add an unresolvable second child at every level.
|
||||
*/
|
||||
void
|
||||
buildInnersDownTo(unsigned int deepest, SHAMapHash childHash, Decoy decoy)
|
||||
{
|
||||
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));
|
||||
|
||||
if (decoy == Decoy::Yes)
|
||||
{
|
||||
// The decoy sits at branch 1, which is free while pathKey is zero, so the
|
||||
// compressed-inner-node parser sees two distinct branches.
|
||||
if (branch == 1)
|
||||
Throw<std::logic_error>("DeepChain: decoy branch collides with real child");
|
||||
|
||||
// Derived from the depth, so each level posts its own read. No node stands behind
|
||||
// this hash, so a read for it stays outstanding.
|
||||
UInt256 decoyHash;
|
||||
decoyHash.begin()[0] = 0xDE;
|
||||
decoyHash.begin()[1] = 0xC0;
|
||||
decoyHash.begin()[2] = static_cast<unsigned char>(depth);
|
||||
s.addBitString(decoyHash);
|
||||
s.add8(1);
|
||||
}
|
||||
|
||||
s.add8(kWireTypeCompressedInner);
|
||||
|
||||
auto node = SHAMapTreeNode::makeFromWire(makeSlice(s.peekData()));
|
||||
childHash = node->getHash();
|
||||
nodes[depth] = std::move(node);
|
||||
}
|
||||
|
||||
rootHash = childHash;
|
||||
}
|
||||
};
|
||||
|
||||
} // namespace xrpl::tests
|
||||
100
src/tests/libxrpl/shamap/InnerNode.h
Normal file
100
src/tests/libxrpl/shamap/InnerNode.h
Normal file
@@ -0,0 +1,100 @@
|
||||
#pragma once
|
||||
|
||||
#include <xrpl/basics/SHAMapHash.h>
|
||||
#include <xrpl/basics/Slice.h>
|
||||
#include <xrpl/basics/base_uint.h>
|
||||
#include <xrpl/basics/contract.h>
|
||||
#include <xrpl/protocol/Serializer.h>
|
||||
#include <xrpl/shamap/SHAMapInnerNode.h>
|
||||
#include <xrpl/shamap/SHAMapTreeNode.h>
|
||||
|
||||
#include <array>
|
||||
#include <cstdint>
|
||||
#include <stdexcept>
|
||||
#include <vector>
|
||||
|
||||
namespace xrpl::tests {
|
||||
|
||||
/**
|
||||
* One child of an inner node a test assembles by hand.
|
||||
*/
|
||||
struct InnerChild
|
||||
{
|
||||
unsigned int branch{};
|
||||
SHAMapHash hash;
|
||||
};
|
||||
|
||||
/**
|
||||
* Assemble an inner node in the wire format's full form.
|
||||
*
|
||||
* A full inner node is all 16 branch hashes back to back in branch order,
|
||||
* followed by the wire type byte. A branch absent from `children`, or given a
|
||||
* zero hash, yields an empty branch, since the parser derives which branches
|
||||
* exist from which hashes are non-zero. The node's hash is already correct,
|
||||
* since makeFromWire() computes it.
|
||||
*
|
||||
* @param children The branch and hash of each child to record. Throws if a
|
||||
* branch is at or above SHAMapInnerNode::kBranchFactor, or if one
|
||||
* branch is named twice.
|
||||
* @return The node, or nullptr if the bytes do not parse.
|
||||
*/
|
||||
[[nodiscard]] inline SHAMapTreeNodePtr
|
||||
makeFullInnerNode(std::vector<InnerChild> const& children)
|
||||
{
|
||||
// The full form carries a slot for every branch, so a child's position comes from the slot it
|
||||
// is written to.
|
||||
std::array<UInt256, SHAMapInnerNode::kBranchFactor> hashes{};
|
||||
|
||||
// A duplicate branch is refused here, and one past the last to match the compressed parser.
|
||||
std::uint32_t seen = 0;
|
||||
for (auto const& child : children)
|
||||
{
|
||||
if (child.branch >= SHAMapInnerNode::kBranchFactor)
|
||||
Throw<std::logic_error>("makeFullInnerNode: branch is past the last one");
|
||||
|
||||
auto const bit = 1u << child.branch;
|
||||
if ((seen & bit) != 0)
|
||||
Throw<std::logic_error>("makeFullInnerNode: branch named twice");
|
||||
seen |= bit;
|
||||
|
||||
hashes.at(child.branch) = child.hash.asUInt256();
|
||||
}
|
||||
|
||||
Serializer s;
|
||||
for (auto const& hash : hashes)
|
||||
s.addBitString(hash);
|
||||
s.add8(kWireTypeInner);
|
||||
|
||||
return SHAMapTreeNode::makeFromWire(makeSlice(s.peekData()));
|
||||
}
|
||||
|
||||
/**
|
||||
* Assemble an inner node in the wire format's compressed form.
|
||||
*
|
||||
* A compressed inner node is one 33-byte chunk per child, each a hash followed
|
||||
* by the branch it sits on, and then the wire type byte. Its hash is already
|
||||
* correct, for the same reason makeFullInnerNode()'s is.
|
||||
*
|
||||
* Nothing is checked here, unlike in makeFullInnerNode(). The branch travels as
|
||||
* one byte, so any value up to 255 reaches the parser, which refuses a branch
|
||||
* at or above SHAMapInnerNode::kBranchFactor and lets a repeated branch
|
||||
* overwrite the hash recorded for it.
|
||||
*
|
||||
* @param children The branch and hash of each child to record.
|
||||
* @return The node, or nullptr if the bytes do not parse.
|
||||
*/
|
||||
[[nodiscard]] inline SHAMapTreeNodePtr
|
||||
makeCompressedInnerNode(std::vector<InnerChild> const& children)
|
||||
{
|
||||
Serializer s;
|
||||
for (auto const& child : children)
|
||||
{
|
||||
s.addBitString(child.hash.asUInt256());
|
||||
s.add8(static_cast<unsigned char>(child.branch));
|
||||
}
|
||||
s.add8(kWireTypeCompressedInner);
|
||||
|
||||
return SHAMapTreeNode::makeFromWire(makeSlice(s.peekData()));
|
||||
}
|
||||
|
||||
} // namespace xrpl::tests
|
||||
76
src/tests/libxrpl/shamap/SHAMapAddNode.cpp
Normal file
76
src/tests/libxrpl/shamap/SHAMapAddNode.cpp
Normal file
@@ -0,0 +1,76 @@
|
||||
#include <xrpl/shamap/SHAMapAddNode.h>
|
||||
|
||||
#include <gtest/gtest.h>
|
||||
|
||||
namespace xrpl::tests {
|
||||
|
||||
// get() is a log format, so its wording is pinned here, once. Every other site reads the same tally
|
||||
// by value, through the count and verdict accessors.
|
||||
TEST(SHAMapAddNode, get_names_every_non_empty_count)
|
||||
{
|
||||
EXPECT_EQ(SHAMapAddNode{}.get(), "no nodes processed");
|
||||
EXPECT_EQ(SHAMapAddNode::useful().get(), "good:1");
|
||||
EXPECT_EQ(SHAMapAddNode::invalid().get(), "bad:1");
|
||||
EXPECT_EQ(SHAMapAddNode::duplicate().get(), "dupe:1");
|
||||
|
||||
// Several of a kind are counted, and the counts are joined in a fixed order with a single
|
||||
// space, whichever order they were recorded in.
|
||||
SHAMapAddNode san;
|
||||
san.incInvalid();
|
||||
san.incUseful();
|
||||
san.incUseful();
|
||||
san.incDuplicate();
|
||||
EXPECT_EQ(san.get(), "good:2 bad:1 dupe:1");
|
||||
|
||||
san.reset();
|
||||
EXPECT_EQ(san.get(), "no nodes processed");
|
||||
}
|
||||
|
||||
// The three counts and the verdicts derived from them.
|
||||
TEST(SHAMapAddNode, counts_and_verdicts_agree)
|
||||
{
|
||||
SHAMapAddNode san;
|
||||
EXPECT_EQ(san.getGood(), 0);
|
||||
EXPECT_EQ(san.getBad(), 0);
|
||||
EXPECT_EQ(san.getDuplicate(), 0);
|
||||
EXPECT_FALSE(san.isInvalid());
|
||||
EXPECT_FALSE(san.isUseful());
|
||||
|
||||
// Good counts what produced a good result, and useful is that count being non-zero.
|
||||
san.incUseful();
|
||||
EXPECT_EQ(san.getGood(), 1);
|
||||
EXPECT_TRUE(san.isUseful());
|
||||
EXPECT_TRUE(san.isGood());
|
||||
|
||||
// A duplicate counts toward good. isUseful() here reflects the incUseful() above.
|
||||
san.incDuplicate();
|
||||
EXPECT_EQ(san.getDuplicate(), 1);
|
||||
EXPECT_FALSE(san.isInvalid());
|
||||
EXPECT_TRUE(san.isGood());
|
||||
|
||||
// Bad is a count, so a batch that carries on past a rejected node reports one per node, which
|
||||
// distinguishes "stopped on the first" from "rejected several".
|
||||
san.incInvalid();
|
||||
EXPECT_EQ(san.getBad(), 1);
|
||||
EXPECT_TRUE(san.isInvalid());
|
||||
EXPECT_TRUE(san.isGood()) << "one bad node among two accepted ones is still a good batch";
|
||||
|
||||
san.incInvalid();
|
||||
san.incInvalid();
|
||||
EXPECT_EQ(san.getGood(), 1);
|
||||
EXPECT_EQ(san.getBad(), 3);
|
||||
EXPECT_EQ(san.getDuplicate(), 1);
|
||||
EXPECT_FALSE(san.isGood()) << "more bad nodes than accepted ones is not a good batch";
|
||||
|
||||
// Adding one verdict to another sums every count.
|
||||
SHAMapAddNode total;
|
||||
total += SHAMapAddNode::useful();
|
||||
total += SHAMapAddNode::invalid();
|
||||
total += SHAMapAddNode::invalid();
|
||||
total += SHAMapAddNode::duplicate();
|
||||
EXPECT_EQ(total.getGood(), 1);
|
||||
EXPECT_EQ(total.getBad(), 2);
|
||||
EXPECT_EQ(total.getDuplicate(), 1);
|
||||
}
|
||||
|
||||
} // namespace xrpl::tests
|
||||
File diff suppressed because it is too large
Load Diff
@@ -14,7 +14,9 @@
|
||||
#include <xrpl/shamap/FullBelowCache.h>
|
||||
#include <xrpl/shamap/TreeNodeCache.h>
|
||||
|
||||
#include <atomic>
|
||||
#include <chrono>
|
||||
#include <cstddef>
|
||||
#include <cstdint>
|
||||
#include <memory>
|
||||
#include <stdexcept>
|
||||
@@ -26,16 +28,27 @@ class TestNodeFamily : public Family
|
||||
private:
|
||||
std::unique_ptr<node_store::Database> db_;
|
||||
|
||||
// Declared before the two caches, which both bind a reference to it in the initializer list.
|
||||
TestStopwatch clock_;
|
||||
|
||||
std::shared_ptr<FullBelowCache> fbCache_;
|
||||
std::shared_ptr<TreeNodeCache> tnCache_;
|
||||
|
||||
TestStopwatch clock_;
|
||||
node_store::DummyScheduler scheduler_;
|
||||
|
||||
beast::Journal const j_;
|
||||
|
||||
// Written from whichever nodestore reader thread reports the miss, so read back atomically.
|
||||
std::atomic<std::size_t> missingBySeqReports_ = 0;
|
||||
std::atomic<std::uint32_t> missingBySeqRefNum_ = 0;
|
||||
|
||||
public:
|
||||
TestNodeFamily(beast::Journal j)
|
||||
/**
|
||||
* @param j The journal to log through.
|
||||
* @param readThreads How many nodestore reader threads to run asynchronous
|
||||
* fetches on.
|
||||
*/
|
||||
explicit TestNodeFamily(beast::Journal j, int readThreads = 1)
|
||||
: fbCache_(std::make_shared<FullBelowCache>("App family full below cache", clock_, j))
|
||||
, tnCache_(
|
||||
std::make_shared<TreeNodeCache>(
|
||||
@@ -50,7 +63,7 @@ public:
|
||||
testSection.set(Keys::kType, "memory");
|
||||
testSection.set(Keys::kPath, "SHAMap_test");
|
||||
db_ = node_store::Manager::instance().makeDatabase(
|
||||
megabytes(4), scheduler_, 1, testSection, j);
|
||||
megabytes(4), scheduler_, readThreads, testSection, j);
|
||||
}
|
||||
|
||||
node_store::Database&
|
||||
@@ -90,14 +103,28 @@ public:
|
||||
tnCache_->sweep();
|
||||
}
|
||||
|
||||
/**
|
||||
* Record the report and throw, standing in for Family's real acquisition
|
||||
* machinery.
|
||||
*
|
||||
* @param refNum Sequence of the ledger with the missing node. Recorded, and
|
||||
* readable through missingBySeqRefNum().
|
||||
* @param nodeHash Hash of the missing node. Unused.
|
||||
*/
|
||||
void
|
||||
missingNodeAcquireBySeq(
|
||||
[[maybe_unused]] std::uint32_t refNum,
|
||||
[[maybe_unused]] UInt256 const& nodeHash) override
|
||||
missingNodeAcquireBySeq(std::uint32_t refNum, [[maybe_unused]] UInt256 const& nodeHash) override
|
||||
{
|
||||
missingBySeqRefNum_.store(refNum, std::memory_order_release);
|
||||
++missingBySeqReports_;
|
||||
Throw<std::runtime_error>("missing node");
|
||||
}
|
||||
|
||||
/**
|
||||
* Throw, standing in for Family's real acquisition machinery. Uncounted.
|
||||
*
|
||||
* @param refHash Hash of the ledger with the missing node. Unused.
|
||||
* @param refNum Sequence of the ledger with the missing node. Unused.
|
||||
*/
|
||||
void
|
||||
missingNodeAcquireByHash(
|
||||
[[maybe_unused]] UInt256 const& refHash,
|
||||
@@ -106,6 +133,29 @@ public:
|
||||
Throw<std::runtime_error>("missing node");
|
||||
}
|
||||
|
||||
/**
|
||||
* How many times a map of this family has withdrawn its claim of being
|
||||
* complete in the database. Counted per family, not per map.
|
||||
*
|
||||
* @return The number of missingNodeAcquireBySeq() calls so far.
|
||||
*/
|
||||
[[nodiscard]] std::size_t
|
||||
missingBySeqReports() const
|
||||
{
|
||||
return missingBySeqReports_.load(std::memory_order_acquire);
|
||||
}
|
||||
|
||||
/**
|
||||
* The ledger sequence the most recent such report named.
|
||||
*
|
||||
* @return The sequence, or zero if nothing has been reported yet.
|
||||
*/
|
||||
[[nodiscard]] std::uint32_t
|
||||
missingBySeqRefNum() const
|
||||
{
|
||||
return missingBySeqRefNum_.load(std::memory_order_acquire);
|
||||
}
|
||||
|
||||
void
|
||||
reset() override
|
||||
{
|
||||
|
||||
@@ -42,7 +42,7 @@ ConsensusTransSetSF::gotNode(
|
||||
|
||||
nodeCache_.insert(nodeHash, nodeData);
|
||||
|
||||
if ((type == SHAMapNodeType::TnTransactionNm) && (nodeData.size() > 16))
|
||||
if ((type == SHAMapNodeType::TnTransactionNm) && (nodeData.size() >= kMinTxNodeBytesToParse))
|
||||
{
|
||||
// this is a transaction, and we didn't have it
|
||||
JLOG(j_.debug()) << "Node on our acquiring TX set is TXN we may not have";
|
||||
|
||||
@@ -9,6 +9,7 @@
|
||||
#include <xrpl/shamap/SHAMapSyncFilter.h>
|
||||
#include <xrpl/shamap/SHAMapTreeNode.h>
|
||||
|
||||
#include <cstddef>
|
||||
#include <cstdint>
|
||||
#include <optional>
|
||||
|
||||
@@ -24,6 +25,15 @@ class ConsensusTransSetSF : public SHAMapSyncFilter
|
||||
public:
|
||||
using NodeCache = TaggedCache<SHAMapHash, Blob>;
|
||||
|
||||
/**
|
||||
* The size a node's hash-prefixed wire data must reach before gotNode()
|
||||
* tries to parse and resubmit it as a transaction. One byte past the
|
||||
* smallest a hash-prefixed SHAMap leaf can be, which a signed transaction
|
||||
* clears.
|
||||
*/
|
||||
static constexpr std::size_t kMinTxNodeBytesToParse =
|
||||
sizeof(std::uint32_t) + kMinShaMapItemBytes + 1;
|
||||
|
||||
ConsensusTransSetSF(Application& app, NodeCache& nodeCache);
|
||||
|
||||
// Note that the nodeData is overwritten by this call
|
||||
|
||||
@@ -30,9 +30,9 @@
|
||||
namespace xrpl {
|
||||
|
||||
// A ledger we are trying to acquire
|
||||
class InboundLedger final : public TimeoutCounter,
|
||||
public std::enable_shared_from_this<InboundLedger>,
|
||||
public CountedObject<InboundLedger>
|
||||
class InboundLedger : public TimeoutCounter,
|
||||
public std::enable_shared_from_this<InboundLedger>,
|
||||
public CountedObject<InboundLedger>
|
||||
{
|
||||
public:
|
||||
using ClockType = beast::AbstractClock<std::chrono::steady_clock>;
|
||||
@@ -44,13 +44,29 @@ public:
|
||||
CONSENSUS // We believe the consensus round requires this ledger
|
||||
};
|
||||
|
||||
/**
|
||||
* How long to wait between retries, and so how long one timeout takes.
|
||||
*/
|
||||
static constexpr std::chrono::milliseconds kRetryInterval{3000};
|
||||
|
||||
/**
|
||||
* @param app The application to run in.
|
||||
* @param hash The ledger to acquire.
|
||||
* @param seq Its sequence, or zero if not known yet.
|
||||
* @param reason Why it is being acquired.
|
||||
* @param clock The clock touch() records against.
|
||||
* @param peerSet Which peers to ask, and how to reach them.
|
||||
* @param retryInterval How long to wait between retries. TimeoutCounter
|
||||
* requires more than 10ms and less than 30s.
|
||||
*/
|
||||
InboundLedger(
|
||||
Application& app,
|
||||
UInt256 const& hash,
|
||||
std::uint32_t seq,
|
||||
Reason reason,
|
||||
ClockType&,
|
||||
std::unique_ptr<PeerSet> peerSet);
|
||||
ClockType& clock,
|
||||
std::unique_ptr<PeerSet> peerSet,
|
||||
std::chrono::milliseconds retryInterval = kRetryInterval);
|
||||
|
||||
~InboundLedger() override;
|
||||
|
||||
@@ -59,7 +75,11 @@ public:
|
||||
update(std::uint32_t seq);
|
||||
|
||||
/**
|
||||
* Returns true if we got all the data.
|
||||
* Whether the acquisition succeeded and its ledger has been settled. Every
|
||||
* path that sets this settles the ledger first, so a caller that sees it
|
||||
* may use the ledger directly.
|
||||
*
|
||||
* @return Whether the ledger is complete and settled.
|
||||
*/
|
||||
bool
|
||||
isComplete() const
|
||||
@@ -68,7 +88,7 @@ public:
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns false if we failed to get the data.
|
||||
* @return Whether the acquisition has failed.
|
||||
*/
|
||||
bool
|
||||
isFailed() const
|
||||
@@ -76,10 +96,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 +148,35 @@ 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, publish its outcome, and signal whatever is
|
||||
* waiting on it. Runs at most once. Call under mtx_, which the flags
|
||||
* written here require. Settles before publishing, since isComplete() is
|
||||
* read without mtx_.
|
||||
*/
|
||||
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();
|
||||
|
||||
@@ -137,11 +186,22 @@ private:
|
||||
void
|
||||
tryDB(node_store::Database& srcDB);
|
||||
|
||||
void
|
||||
done();
|
||||
/**
|
||||
* Whether either map of the ledger being acquired has been found invalid.
|
||||
*
|
||||
* A walk returns a bare list of hashes, so this is what tells a satisfied
|
||||
* map from an abandoned one. Callers ask it before reading emptiness as
|
||||
* "nothing left to fetch". See SHAMap::addKnownNode for why the verdict is
|
||||
* final.
|
||||
*
|
||||
* @return Whether either map is Invalid, and false while there is no ledger
|
||||
* yet, since then there is no map to judge.
|
||||
*/
|
||||
[[nodiscard]] bool
|
||||
hasInvalidMap() const;
|
||||
|
||||
void
|
||||
onTimer(bool progress, ScopedLockType& peerSetLock) override;
|
||||
onTimer(bool progress, ScopedLockType& sl) override;
|
||||
|
||||
std::size_t
|
||||
getPeerCount() const;
|
||||
@@ -155,6 +215,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,
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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 both
|
||||
// maps are still sound here.
|
||||
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;
|
||||
@@ -222,6 +236,12 @@ InboundLedger::neededStateHashes(int max, SHAMapSyncFilter const* filter) const
|
||||
return neededHashes(ledger_->header().accountHash, ledger_->stateMap(), max, filter);
|
||||
}
|
||||
|
||||
bool
|
||||
InboundLedger::hasInvalidMap() const
|
||||
{
|
||||
return ledger_ && !ledger_->mapsValid();
|
||||
}
|
||||
|
||||
// See how much of the ledger data is stored locally
|
||||
// Data found in a fetch pack will be stored
|
||||
void
|
||||
@@ -241,7 +261,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 +334,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))
|
||||
{
|
||||
@@ -326,14 +345,32 @@ InboundLedger::tryDB(node_store::Database& srcDB)
|
||||
}
|
||||
}
|
||||
|
||||
// Judged here rather than at the setImmutable() below, which runs only once both flags are
|
||||
// set: one map can be abandoned while the other is merely incomplete.
|
||||
if (hasInvalidMap())
|
||||
{
|
||||
JLOG(journal_.warn()) << "Ledger " << hash_ << " found locally has an invalid map";
|
||||
failed_ = true;
|
||||
return;
|
||||
}
|
||||
|
||||
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 sees a ledger
|
||||
// this function has finished with. Reached despite the guard above because trigger()
|
||||
// walks the state map with mtx_ released, so that walk can reach the verdict in between.
|
||||
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 +456,46 @@ InboundLedger::done()
|
||||
signaled_ = true;
|
||||
touch();
|
||||
|
||||
// Settled here, and complete_ published only once it is settled. isComplete() is read without
|
||||
// mtx_, by InboundLedgers::acquire() among others, and LedgerHistory::insert() and
|
||||
// LedgerHolder::set() each require an immutable ledger. tryDB() already settles its own
|
||||
// result and sets complete_ itself, so that path arrives here with the ledger immutable and
|
||||
// only the reporting below left to do.
|
||||
bool const haveEverything = haveHeader_ && haveState_ && haveTransactions_;
|
||||
if (!failed_ && ledger_ && (complete_ || haveEverything))
|
||||
{
|
||||
XRPL_ASSERT(
|
||||
ledger_->header().seq < kXrpLedgerEarliestFees || ledger_->read(keylet::feeSettings()),
|
||||
"xrpl::InboundLedger::done : valid ledger fees");
|
||||
// trigger() walks the state map with mtx_ released, so that walk can reach the verdict
|
||||
// after the flags said there was nothing left to fetch. A race rather than a broken
|
||||
// invariant, and one that peer data produces, so this recovers rather than asserts.
|
||||
// setInvalid() outranks Immutable, so a walk that reaches the verdict after both maps
|
||||
// have been settled leaves an immutable ledger with an invalid map.
|
||||
SOMETIMES(hasInvalidMap(), "xrpl::InboundLedger::done : map invalidated by a race");
|
||||
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_, sees the claim withdrawn.
|
||||
complete_ = false;
|
||||
failed_ = true;
|
||||
}
|
||||
else
|
||||
{
|
||||
complete_ = true;
|
||||
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,28 +504,13 @@ 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_)
|
||||
// Read through getLedger(), which consults failed_ itself, so failure and a null ledger_
|
||||
// are both refused by one check. checkAccept() requires a non-null ledger.
|
||||
if (auto const ledger = self->getLedger(); self->complete_ && ledger)
|
||||
{
|
||||
self->app_.getLedgerMaster().checkAccept(self->getLedger());
|
||||
self->app_.getLedgerMaster().checkAccept(ledger);
|
||||
self->app_.getLedgerMaster().tryAdvance();
|
||||
}
|
||||
else
|
||||
@@ -497,6 +559,7 @@ InboundLedger::trigger(std::shared_ptr<Peer> const& peer, TriggerReason reason)
|
||||
if (failed_)
|
||||
{
|
||||
JLOG(journal_.warn()) << " failed local for " << hash_;
|
||||
done();
|
||||
return;
|
||||
}
|
||||
}
|
||||
@@ -513,7 +576,20 @@ InboundLedger::trigger(std::shared_ptr<Peer> const& peer, TriggerReason reason)
|
||||
{
|
||||
auto need = getNeededHashes();
|
||||
|
||||
if (!need.empty())
|
||||
// The validity check runs ahead of the emptiness test, since getNeededHashes() walks
|
||||
// both maps and can reach the verdict itself. The result is read once, so the hint
|
||||
// below and the test it feeds share one observation of a map another thread can be
|
||||
// invalidating. The claim is withdrawn alongside the failure, so failed_ and
|
||||
// complete_ never both read true.
|
||||
bool const invalidMap = hasInvalidMap();
|
||||
SOMETIMES(invalidMap, "xrpl::InboundLedger::trigger : map is invalid");
|
||||
if (invalidMap)
|
||||
{
|
||||
JLOG(journal_.warn()) << "Acquire " << hash_ << " has an invalid map";
|
||||
failed_ = true;
|
||||
complete_ = false;
|
||||
}
|
||||
else if (!need.empty())
|
||||
{
|
||||
protocol::TMGetObjectByHash tmBH;
|
||||
bool typeSet = false;
|
||||
@@ -550,11 +626,11 @@ InboundLedger::trigger(std::shared_ptr<Peer> const& peer, TriggerReason reason)
|
||||
}
|
||||
else
|
||||
{
|
||||
// The tail of this function settles the ledger and reports it complete.
|
||||
JLOG(journal_.info()) << "getNeededHashes says acquire is complete";
|
||||
haveHeader_ = true;
|
||||
haveTransactions_ = true;
|
||||
haveState_ = true;
|
||||
complete_ = true;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -617,27 +693,31 @@ InboundLedger::trigger(std::shared_ptr<Peer> const& peer, TriggerReason reason)
|
||||
{
|
||||
AccountStateSF filter(ledger_->stateMap().family().db(), app_.getLedgerMaster());
|
||||
|
||||
// Release the lock while we process the large state map
|
||||
// Release the lock while the state map is walked. mtx_ is recursive, so an sl.unlock()
|
||||
// under onTimer() or the addPeers() callback leaves mtx_ held. The flags are re-read
|
||||
// below because another packet can be handled while this walk runs.
|
||||
sl.unlock();
|
||||
auto nodes = ledger_->stateMap().getMissingNodes(kMissingNodesFind, &filter);
|
||||
sl.lock();
|
||||
|
||||
// The validity check runs outside the flags' guard below, since the verdict is about
|
||||
// the map rather than this round: it holds even if another thread reported the ledger
|
||||
// complete while the lock was released.
|
||||
bool const walkAbandonedMap = hasInvalidMap();
|
||||
SOMETIMES(walkAbandonedMap, "xrpl::InboundLedger::trigger : map abandoned by its walk");
|
||||
if (walkAbandonedMap)
|
||||
{
|
||||
JLOG(journal_.warn()) << "Ledger " << hash_ << " has a map its walk abandoned";
|
||||
failed_ = true;
|
||||
complete_ = false;
|
||||
}
|
||||
// Make sure nothing happened while we released the lock
|
||||
if (!failed_ && !complete_ && !haveState_)
|
||||
else if (!failed_ && !complete_ && !haveState_)
|
||||
{
|
||||
if (nodes.empty())
|
||||
{
|
||||
if (!ledger_->stateMap().isValid())
|
||||
{
|
||||
failed_ = true;
|
||||
}
|
||||
else
|
||||
{
|
||||
haveState_ = true;
|
||||
|
||||
if (haveTransactions_)
|
||||
complete_ = true;
|
||||
}
|
||||
// The test above already caught a map the walk abandoned, so this one is sound.
|
||||
haveState_ = true;
|
||||
}
|
||||
else
|
||||
{
|
||||
@@ -699,9 +779,6 @@ InboundLedger::trigger(std::shared_ptr<Peer> const& peer, TriggerReason reason)
|
||||
else
|
||||
{
|
||||
haveTransactions_ = true;
|
||||
|
||||
if (haveState_)
|
||||
complete_ = true;
|
||||
}
|
||||
}
|
||||
else
|
||||
@@ -726,11 +803,14 @@ InboundLedger::trigger(std::shared_ptr<Peer> const& peer, TriggerReason reason)
|
||||
}
|
||||
}
|
||||
|
||||
if (complete_ || failed_)
|
||||
// Having every part is not yet a completed acquisition: done() settles the ledger first and
|
||||
// only then publishes complete_. Called with mtx_ still held, as done() documents, so the flags
|
||||
// it writes are not written unlocked; mtx_ is recursive, so a caller that already holds it is
|
||||
// unaffected.
|
||||
if (failed_ || (haveHeader_ && haveState_ && haveTransactions_))
|
||||
{
|
||||
JLOG(journal_.debug()) << "Done:" << (complete_ ? " complete" : "")
|
||||
<< (failed_ ? " failed " : " ") << ledger_->header().seq;
|
||||
sl.unlock();
|
||||
JLOG(journal_.debug()) << "Done:" << (failed_ ? " failed " : " have everything ")
|
||||
<< ledger_->header().seq;
|
||||
done();
|
||||
}
|
||||
}
|
||||
@@ -774,15 +854,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 +894,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 +914,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();
|
||||
|
||||
@@ -928,11 +1027,10 @@ InboundLedger::receiveNode(
|
||||
haveState_ = true;
|
||||
}
|
||||
|
||||
// done() settles the ledger before publishing complete_, so having every part is reported
|
||||
// there rather than here.
|
||||
if (haveTransactions_ && haveState_)
|
||||
{
|
||||
complete_ = true;
|
||||
done();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1099,6 +1197,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();
|
||||
}
|
||||
|
||||
|
||||
@@ -168,15 +168,8 @@ public:
|
||||
data.emplace_back(*nodeID, std::move(treeNode));
|
||||
}
|
||||
|
||||
auto const san = ta->takeNodes(std::move(data), peer);
|
||||
if (san.isInvalid())
|
||||
{
|
||||
peer->charge(resource::kFeeInvalidData, "ledger_data invalid");
|
||||
}
|
||||
else if (!san.isUseful())
|
||||
{
|
||||
peer->charge(resource::kFeeUselessData, "ledger_data useless");
|
||||
}
|
||||
// takeNodes() charges the peer and records the batch itself, so the verdict is discarded.
|
||||
static_cast<void>(ta->takeNodes(std::move(data), peer));
|
||||
}
|
||||
|
||||
void
|
||||
|
||||
@@ -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);
|
||||
|
||||
|
||||
@@ -8,6 +8,7 @@
|
||||
|
||||
#include <boost/asio/basic_waitable_timer.hpp>
|
||||
|
||||
#include <atomic>
|
||||
#include <chrono>
|
||||
#include <cstdint>
|
||||
#include <memory>
|
||||
@@ -50,6 +51,9 @@ namespace xrpl {
|
||||
* whether to postpone failure and reset the timeout. However, if it can
|
||||
* complete all its work in one synchronous step (while it holds the lock), then
|
||||
* it can ignore `progress_`.
|
||||
*
|
||||
* `isDone` is not terminal for every subtype: TransactionAcquire::stillNeed()
|
||||
* clears `failed_` and calls `setTimer` again.
|
||||
*/
|
||||
class TimeoutCounter
|
||||
{
|
||||
@@ -98,9 +102,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.
|
||||
@@ -126,10 +135,27 @@ protected:
|
||||
*/
|
||||
UInt256 const hash_;
|
||||
int timeouts_{0};
|
||||
bool complete_{false};
|
||||
bool failed_{false};
|
||||
|
||||
// complete_ and failed_ are read without mtx_, so they are atomic rather than guarded: a
|
||||
// reader that sees either flag set also sees the work the writer did before setting it.
|
||||
static_assert(std::atomic<bool>::is_always_lock_free);
|
||||
|
||||
/**
|
||||
* Whether forward progress has been made.
|
||||
* Whether the task finished successfully. Each subclass sets it only once
|
||||
* the work it reports on is finished, so a reader may act on it without
|
||||
* taking mtx_.
|
||||
*/
|
||||
std::atomic<bool> complete_{false};
|
||||
|
||||
/**
|
||||
* Whether the task gave up.
|
||||
*/
|
||||
std::atomic<bool> failed_{false};
|
||||
|
||||
/**
|
||||
* Whether forward progress has been made since invokeOnTimer() last ran.
|
||||
* Each subtype defines what counts as progress, and may read this for
|
||||
* decisions of its own.
|
||||
*/
|
||||
bool progress_{false};
|
||||
/**
|
||||
|
||||
@@ -8,7 +8,9 @@
|
||||
|
||||
#include <xrpl/basics/Log.h>
|
||||
#include <xrpl/basics/base_uint.h>
|
||||
#include <xrpl/beast/utility/instrumentation.h>
|
||||
#include <xrpl/core/Job.h>
|
||||
#include <xrpl/resource/Fees.h>
|
||||
#include <xrpl/server/NetworkOPs.h>
|
||||
#include <xrpl/shamap/SHAMap.h>
|
||||
#include <xrpl/shamap/SHAMapAddNode.h>
|
||||
@@ -18,6 +20,7 @@
|
||||
#include <xrpl.pb.h>
|
||||
|
||||
#include <algorithm>
|
||||
#include <chrono>
|
||||
#include <cstddef>
|
||||
#include <exception>
|
||||
#include <memory>
|
||||
@@ -26,22 +29,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 +52,28 @@ 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)
|
||||
{
|
||||
@@ -100,6 +111,25 @@ TransactionAcquire::pmDowncast()
|
||||
return shared_from_this();
|
||||
}
|
||||
|
||||
void
|
||||
TransactionAcquire::recordAsked(std::shared_ptr<Peer> const& peer)
|
||||
{
|
||||
// Each request renews the pass chargeLateReply() reads, so a peer answering the request
|
||||
// just sent to it is free whatever it answered in an earlier round.
|
||||
if (peer)
|
||||
{
|
||||
requestedPeers_.insert(peer->id());
|
||||
lateReplyGranted_.erase(peer->id());
|
||||
return;
|
||||
}
|
||||
|
||||
// A broadcast goes to every peer the set tracks, so each of them has been asked.
|
||||
auto const& ids = peerSet_->getPeerIds();
|
||||
requestedPeers_.insert(ids.begin(), ids.end());
|
||||
for (auto const id : ids)
|
||||
lateReplyGranted_.erase(id);
|
||||
}
|
||||
|
||||
void
|
||||
TransactionAcquire::trigger(std::shared_ptr<Peer> const& peer)
|
||||
{
|
||||
@@ -127,6 +157,7 @@ TransactionAcquire::trigger(std::shared_ptr<Peer> const& peer)
|
||||
tmGL.set_querytype(protocol::qtINDIRECT);
|
||||
|
||||
*(tmGL.add_nodeids()) = SHAMapNodeID().getRawString();
|
||||
recordAsked(peer);
|
||||
peerSet_->sendRequest(tmGL, peer);
|
||||
}
|
||||
else if (!map_->isValid())
|
||||
@@ -165,6 +196,7 @@ TransactionAcquire::trigger(std::shared_ptr<Peer> const& peer)
|
||||
{
|
||||
*tmGL.add_nodeids() = node.first.getRawString();
|
||||
}
|
||||
recordAsked(peer);
|
||||
peerSet_->sendRequest(tmGL, peer);
|
||||
}
|
||||
}
|
||||
@@ -174,24 +206,59 @@ TransactionAcquire::takeNodes(
|
||||
std::vector<std::pair<SHAMapNodeID, SHAMapTreeNodePtr>> data,
|
||||
std::shared_ptr<Peer> const& peer)
|
||||
{
|
||||
ScopedLockType const sl(mtx_);
|
||||
ScopedLockType sl(mtx_);
|
||||
|
||||
if (complete_)
|
||||
// Read before the call below, which can settle the set itself.
|
||||
bool const wasSettled = isDone();
|
||||
|
||||
auto const san = takeNodesLocked(std::move(data), peer, sl);
|
||||
|
||||
// A batch that advanced the map must keep the next timer tick from counting a timeout against
|
||||
// it. A duplicate counts as an answer: an honest second responder to trigger()'s fan-out has
|
||||
// replied, so no timeout is owed. A reply to a set already settled on entry owes nothing,
|
||||
// since the allowance it spends belongs to the round that ended.
|
||||
if (!wasSettled && (san.isUseful() || san.getDuplicate() > 0))
|
||||
progress_ = true;
|
||||
|
||||
return san;
|
||||
}
|
||||
|
||||
SHAMapAddNode
|
||||
TransactionAcquire::takeNodesLocked(
|
||||
std::vector<std::pair<SHAMapNodeID, SHAMapTreeNodePtr>> data,
|
||||
std::shared_ptr<Peer> const& peer,
|
||||
ScopedLockType& sl)
|
||||
{
|
||||
// A reply that arrives after the set is settled - by completing it, or by a different packet
|
||||
// failing it. trigger() sends to every peer it was given, so any of their replies, including
|
||||
// another packet from the same peer whose data failed the set, can already be in flight and
|
||||
// could not have known the outcome. Those are solicited, and free: one per peer we asked,
|
||||
// which is what bounds the honest case.
|
||||
//
|
||||
// Past that bound, further data for this hash is a replay - a resend of data already
|
||||
// accepted or now known worthless, not a first-time reply - and serving it is not free
|
||||
// work, so it is charged.
|
||||
if (isDone())
|
||||
{
|
||||
JLOG(journal_.trace()) << "TX set complete";
|
||||
return SHAMapAddNode();
|
||||
JLOG(journal_.trace()) << (complete_ ? "TX set complete" : "TX set failed");
|
||||
|
||||
chargeLateReply(peer, sl);
|
||||
|
||||
return SHAMapAddNode::duplicate();
|
||||
}
|
||||
|
||||
if (failed_)
|
||||
{
|
||||
JLOG(journal_.trace()) << "TX set failed";
|
||||
return SHAMapAddNode();
|
||||
}
|
||||
// Accumulated across the batch, so a packet ending in one bad node still counts the nodes
|
||||
// hooked in ahead of it, as InboundLedger::receiveNode() already does.
|
||||
SHAMapAddNode san;
|
||||
|
||||
try
|
||||
{
|
||||
if (data.empty())
|
||||
{
|
||||
// Defensive: PeerImp rejects an empty node list before dispatch.
|
||||
peer->charge(resource::kFeeInvalidData, "tx_set empty");
|
||||
return SHAMapAddNode::invalid();
|
||||
}
|
||||
|
||||
ConsensusTransSetSF sf(app_, app_.getTempNodeCache());
|
||||
|
||||
@@ -202,36 +269,54 @@ TransactionAcquire::takeNodes(
|
||||
if (haveRoot_)
|
||||
{
|
||||
JLOG(journal_.debug()) << "Got root TXS node, already have it";
|
||||
san.incDuplicate();
|
||||
continue;
|
||||
}
|
||||
else if (!map_->addRootNode(SHAMapHash{hash_}, std::move(d.second), nullptr)
|
||||
.isGood())
|
||||
|
||||
auto const result =
|
||||
map_->addRootNode(SHAMapHash{hash_}, std::move(d.second), nullptr);
|
||||
san += result;
|
||||
|
||||
if (!result.isGood())
|
||||
{
|
||||
JLOG(journal_.warn()) << "TX acquire got bad root node for TX set " << hash_
|
||||
<< " from peer " << peer->id();
|
||||
return SHAMapAddNode::invalid();
|
||||
}
|
||||
else
|
||||
{
|
||||
haveRoot_ = true;
|
||||
// addRootNode only rejects a hash mismatch, so the timer will retry with
|
||||
// another peer.
|
||||
peer->charge(resource::kFeeInvalidData, "tx_set root hash mismatch");
|
||||
return san;
|
||||
}
|
||||
|
||||
haveRoot_ = true;
|
||||
continue;
|
||||
}
|
||||
else if (!map_->addKnownNode(d.first, std::move(d.second), &sf).isGood())
|
||||
|
||||
auto const result = map_->addKnownNode(d.first, std::move(d.second), &sf);
|
||||
san += result;
|
||||
|
||||
if (!result.isGood())
|
||||
{
|
||||
JLOG(journal_.warn()) << "TX acquire got bad non-root node " << d.first
|
||||
<< " for TX set " << hash_ << " from peer " << peer->id();
|
||||
return SHAMapAddNode::invalid();
|
||||
// A bad node leaves the map sound, so leave that retry to the timer rather than
|
||||
// re-requesting from the peer that just sent us bad data.
|
||||
peer->charge(resource::kFeeInvalidData, "tx_set node invalid");
|
||||
return san;
|
||||
}
|
||||
}
|
||||
|
||||
trigger(peer);
|
||||
progress_ = true;
|
||||
return SHAMapAddNode::useful();
|
||||
return san;
|
||||
}
|
||||
catch (std::exception const& ex)
|
||||
{
|
||||
JLOG(journal_.error()) << "Peer " << peer->id()
|
||||
<< " sent us junky transaction node data: " << ex.what();
|
||||
return SHAMapAddNode::invalid();
|
||||
JLOG(journal_.error()) << "TX acquire threw while taking nodes for TX set " << hash_
|
||||
<< " from peer " << peer->id() << ": " << ex.what();
|
||||
// Whatever the batch hooked in before the throw stands, so the tally it reached is what
|
||||
// the caller is told. The timer owns the retry. The sender keeps its fee, since the
|
||||
// classes that reach here are raised by this code rather than by the data.
|
||||
san.incInvalid();
|
||||
return san;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -254,13 +339,36 @@ TransactionAcquire::init(int numPeers)
|
||||
setTimer(sl);
|
||||
}
|
||||
|
||||
void
|
||||
TransactionAcquire::chargeLateReply(std::shared_ptr<Peer> const& peer, ScopedLockType&)
|
||||
{
|
||||
if (!requestedPeers_.contains(peer->id()) || !lateReplyGranted_.insert(peer->id()).second)
|
||||
peer->charge(resource::kFeeUselessData, "tx_set data after the set was settled");
|
||||
}
|
||||
|
||||
void
|
||||
TransactionAcquire::stillNeed()
|
||||
{
|
||||
ScopedLockType const sl(mtx_);
|
||||
ScopedLockType sl(mtx_);
|
||||
|
||||
timeouts_ = std::min<int>(timeouts_, kNormTimeouts);
|
||||
|
||||
// A running acquisition keeps the wait it has, rather than restarting it for every consensus
|
||||
// round that asks for the set again.
|
||||
if (!failed_)
|
||||
return;
|
||||
|
||||
failed_ = false;
|
||||
|
||||
// lateReplyGranted_ is left alone. The free allowance is earned one request at a time, so a
|
||||
// peer keeps the pass it already spent until recordAsked() records another request to it. The
|
||||
// timer restarted below is what sends those requests.
|
||||
|
||||
// Restarting the timer is what resumes the acquisition. expires_after() cancels whatever wait
|
||||
// was outstanding, so the timer holds at most one wait at a time. A job queueJob() already
|
||||
// handed to the JobQueue is not canceled by that and still runs one invokeOnTimer(), which
|
||||
// re-arms this same timer and so folds back into the one chain.
|
||||
setTimer(sl);
|
||||
}
|
||||
|
||||
} // namespace xrpl
|
||||
|
||||
@@ -12,8 +12,10 @@
|
||||
#include <xrpl/shamap/SHAMapNodeID.h>
|
||||
#include <xrpl/shamap/SHAMapTreeNode.h>
|
||||
|
||||
#include <chrono>
|
||||
#include <cstddef>
|
||||
#include <memory>
|
||||
#include <set>
|
||||
#include <utility>
|
||||
#include <vector>
|
||||
|
||||
@@ -21,16 +23,46 @@ 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;
|
||||
|
||||
/**
|
||||
* Add nodes a peer sent us to the set we are acquiring.
|
||||
*
|
||||
* Charges the peer for data it declines, since only this function holds the
|
||||
* lock that decides the tier. A late reply is bounded by a per-peer
|
||||
* allowance; see lateReplyGranted_.
|
||||
*
|
||||
* @param data The nodes to add, each with its claimed position.
|
||||
* @param peer The peer that sent them, charged here if the data is
|
||||
* declined.
|
||||
* @return The tally of useful, duplicate, and bad nodes in the batch.
|
||||
* Useful and bad can both be nonzero, since only the node the
|
||||
* batch stops on is bad.
|
||||
*/
|
||||
SHAMapAddNode
|
||||
takeNodes(
|
||||
std::vector<std::pair<SHAMapNodeID, SHAMapTreeNodePtr>> data,
|
||||
@@ -39,23 +71,114 @@ public:
|
||||
void
|
||||
init(int startPeers);
|
||||
|
||||
/**
|
||||
* Resume a timed-out acquisition, or leave a running one alone.
|
||||
*
|
||||
* Always clamps the timeout count. An acquisition that failed has its timer
|
||||
* chain stopped, so this also clears the failed flag and restarts the timer;
|
||||
* one that is still running already has a timer pending.
|
||||
*/
|
||||
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};
|
||||
|
||||
/**
|
||||
* Every peer a request has actually been sent to.
|
||||
*
|
||||
* Holds the peers addPeers() selected and trigger() then built a request
|
||||
* for, the unsolicited senders takeNodesLocked() trigger()s directly
|
||||
* outside addPeers(), and the whole of peerSet_->getPeerIds() once a
|
||||
* broadcast has gone out. Selection alone does not enroll a peer: trigger()
|
||||
* records it at the two branches that build a request, so one reached while
|
||||
* the map is invalid, while nothing is missing, or after the acquisition
|
||||
* settled is left out. Recording every peer a request went to, however it
|
||||
* was chosen, bounds the free allowance to peers with a reply in flight.
|
||||
* stillNeed() keeps this set, since a peer already asked stays asked
|
||||
* whichever round its reply arrives in.
|
||||
*/
|
||||
std::set<Peer::ID> requestedPeers_;
|
||||
|
||||
/**
|
||||
* Peers in requestedPeers_ that have already spent the free late reply
|
||||
* their last request earned.
|
||||
*
|
||||
* Membership rather than a count, so each peer's pass is its own: a
|
||||
* peer already in requestedPeers_ is granted one free late reply the
|
||||
* first time it reaches takeNodesLocked()'s isDone() branch, and is
|
||||
* charged on every later one. recordAsked() renews a peer's pass wherever
|
||||
* it records a request, for a targeted request and for a broadcast alike.
|
||||
* So each peer holds one unspent pass, renewed by each request sent to it,
|
||||
* and overlapping requests to one peer share one pass. A revival on its
|
||||
* own renews nothing.
|
||||
*/
|
||||
std::set<Peer::ID> lateReplyGranted_;
|
||||
|
||||
std::unique_ptr<PeerSet> peerSet_;
|
||||
|
||||
void
|
||||
onTimer(bool progress, ScopedLockType& peerSetLock) override;
|
||||
/**
|
||||
* Add nodes a peer sent us, on the lock takeNodes() holds.
|
||||
*
|
||||
* Split out so recording what the batch achieved happens on one exit,
|
||||
* covering the paths that stop the batch early as well.
|
||||
*
|
||||
* @param data The nodes to add, each with its claimed position.
|
||||
* @param peer The peer that sent them, charged here if the data is
|
||||
* declined.
|
||||
* @return The tally of useful, duplicate, and bad nodes in the batch.
|
||||
*/
|
||||
SHAMapAddNode
|
||||
takeNodesLocked(
|
||||
std::vector<std::pair<SHAMapNodeID, SHAMapTreeNodePtr>> data,
|
||||
std::shared_ptr<Peer> const& peer,
|
||||
ScopedLockType& sl);
|
||||
|
||||
/**
|
||||
* Spend this peer's one free late reply, or charge it for replaying.
|
||||
*
|
||||
* Called from takeNodesLocked(), which is where a late reply is
|
||||
* recognized.
|
||||
*
|
||||
* @param peer The peer that sent the reply.
|
||||
* @param sl Proof mtx_ is held, which the allowance sets require.
|
||||
*/
|
||||
void
|
||||
chargeLateReply(std::shared_ptr<Peer> const& peer, ScopedLockType& sl);
|
||||
|
||||
void
|
||||
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();
|
||||
|
||||
void
|
||||
addPeers(std::size_t limit);
|
||||
|
||||
/**
|
||||
* Record every peer a request has just gone to, and renew each one's free
|
||||
* late reply. Call under mtx_.
|
||||
*
|
||||
* The renewal gives each peer one unspent pass, renewed by each request
|
||||
* sent to it, so overlapping requests to one peer share one pass. See
|
||||
* lateReplyGranted_.
|
||||
*
|
||||
* @param peer The peer a targeted request went to, or nullptr for a
|
||||
* broadcast, which reaches every peer the set tracks.
|
||||
*/
|
||||
void
|
||||
recordAsked(std::shared_ptr<Peer> const& peer);
|
||||
|
||||
void
|
||||
trigger(std::shared_ptr<Peer> const&);
|
||||
std::weak_ptr<TimeoutCounter>
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user