Compare commits

...

3 Commits

Author SHA1 Message Date
Bart
a790fb3527 refactor: Add a reusable peer harness for acquisition tests
DeepChain builds the node chains both acquisition suites need: chains
that run to SHAMap::kLeafDepth, which no valid tree holds, and toLeaf()
chains that complete an acquisition. AcquireTestHelpers.h adds
ChargeRecordingPeer, RequestCountingPeerSet, packetFor(), waitFor() and
tallyIs(), so each suite drives an acquisition through the real
gotData() dispatch rather than reproducing it.

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

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

View File

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

View File

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

View File

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

View File

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

View File

@@ -0,0 +1,300 @@
#include <test/app/AcquireTestHelpers.h>
#include <test/jtx/Env.h>
#include <xrpld/app/ledger/InboundLedger.h>
#include <xrpld/app/ledger/InboundLedgers.h>
#include <xrpld/app/ledger/LedgerMaster.h>
#include <xrpl/basics/base_uint.h>
#include <xrpl/basics/chrono.h>
#include <xrpl/beast/unit_test/suite.h>
#include <xrpl/nodestore/NodeObject.h>
#include <xrpl/protocol/HashPrefix.h>
#include <xrpl/protocol/LedgerHeader.h>
#include <xrpl/protocol/Serializer.h>
#include <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);
}
};
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);
}
}
/**
* A ledger whose maps all resolve locally finishes on the spot, and the
* finished ledger is immutable and handed on.
*
* Both entry points into tryDB() are driven: InboundLedgers::acquire(), the
* only caller of init(), and checkLocal(), which reaches done().
*
* @param env The environment to run in.
*/
void
testLocalLedgerCompletesAcquire(jtx::Env& env)
{
testcase("A ledger found locally completes the acquire");
// A chain ending in a real leaf, so the state map is complete rather than merely rooted.
// No transactions, so only the state map is in play.
auto const chain = DeepChain::toLeaf(2, nextSeed());
auto const header = makeHeader(chain);
storeHeader(env, header);
storeStateNodes(env, header, chain, chain.deepestDepth);
// acquire() returns the ledger only once the acquisition is complete and unfailed, so a
// non-null result shows tryDB() found it in the store.
auto const acquired = env.app().getInboundLedgers().acquire(
header.hash, header.seq, InboundLedger::Reason::GENERIC);
BEAST_EXPECT(acquired != nullptr);
if (acquired)
{
BEAST_EXPECT(acquired->isImmutable());
BEAST_EXPECT(acquired->header().hash == header.hash);
}
// init() hands a ledger it completed to LedgerMaster itself, which is what makes it
// available to everything else.
BEAST_EXPECT(env.app().getLedgerMaster().getLedgerByHash(header.hash) != nullptr);
// The failure list stays clear, which is the other arm of done().
BEAST_EXPECT(!env.app().getInboundLedgers().isFailure(header.hash));
// The same ledger through checkLocal(), which unlike init() reaches done(). Everything it
// needs is still in the store, since the first acquisition read rather than consumed it.
auto again = std::make_shared<InboundLedger>(
env.app(),
header.hash,
header.seq,
InboundLedger::Reason::GENERIC,
stopwatch(),
std::make_unique<RequestCountingPeerSet>());
// True because the acquisition ended, which it reports only after done() has run.
BEAST_EXPECT(again->checkLocal());
BEAST_EXPECT(again->isComplete());
BEAST_EXPECT(!again->isFailed());
auto const settled = again->getLedger();
BEAST_EXPECT(settled != nullptr);
if (settled)
BEAST_EXPECT(settled->isImmutable());
BEAST_EXPECT(!env.app().getInboundLedgers().isFailure(header.hash));
}
/**
* The retry timer re-asks, then gives up and signals.
*
* The only case that drives onTimer() rather than trigger() directly,
* which is what covers the give-up: past kLedgerTimeoutRetriesMax the
* acquisition fails itself and done() records that, so the same
* doomed ledger is not asked for again on the next round. It is also
* what the retry interval is a constructor parameter for, since the
* chain runs past kLedgerTimeoutRetriesMax ticks of three seconds
* apiece in production.
*
* The store is empty for this hash and no reply arrives, so every tick
* counts a timeout and the count climbs to the limit. A hash of its own,
* so this case owns its entry in the failure list.
*
* @param env The environment to run in.
*/
void
testTimerRetriesThenGivesUp(jtx::Env& env)
{
testcase("The retry timer re-asks, then gives up");
UInt256 const kUnknownLedger{8};
// One candidate, which onTimer() re-offers on every tick.
auto const candidate = std::make_shared<ChargeRecordingPeer>();
auto peerSet =
std::make_unique<RequestCountingPeerSet>(std::vector<std::shared_ptr<Peer>>{candidate});
auto* const peerSetPtr = peerSet.get();
auto acquire = std::make_shared<TestableInboundLedger>(
env.app(),
kUnknownLedger,
0,
InboundLedger::Reason::GENERIC,
stopwatch(),
std::move(peerSet),
kFastRetry);
BEAST_EXPECT(!env.app().getInboundLedgers().isFailure(kUnknownLedger));
// init() finds the hash absent from the store, so it asks the candidate and arms the retry
// timer. The case counts from here, past those first requests.
acquire->startAcquire();
BEAST_EXPECT(!acquire->isFailed());
int const requestsFromInit = peerSetPtr->requests();
BEAST_EXPECT(requestsFromInit > 0);
BEAST_EXPECT(peerSetPtr->addedPeers() == std::set<Peer::ID>{candidate->id()});
// Every tick asks again, and past kLedgerTimeoutRetriesMax (6) the chain gives up.
BEAST_EXPECT(waitFor([&] { return acquire->isFailed(); }));
BEAST_EXPECT(!acquire->isComplete());
BEAST_EXPECT(peerSetPtr->requests() > requestsFromInit);
// done() remembered the hash, which is what stops the next round asking again.
BEAST_EXPECT(
waitFor([&] { return env.app().getInboundLedgers().isFailure(kUnknownLedger); }));
}
void
run() override
{
// One Env for the suite, since building one costs far more than any case here. Safe
// because every chain is seeded through nextSeed(): the node store, the fetch packs and
// the remembered failures are all shared, and all three are keyed by hash.
jtx::Env env{*this};
testLocalLedgerCompletesAcquire(env);
// Last: the only case that waits out a whole timeout chain.
testTimerRetriesThenGivesUp(env);
}
private:
unsigned int seed_{0};
};
BEAST_DEFINE_TESTSUITE(InboundLedger, app, xrpl);
} // namespace xrpl::test

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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