xrpld
Loading...
Searching...
No Matches
SkipListAcquire.cpp
1#include <xrpld/app/ledger/detail/SkipListAcquire.h>
2
3#include <xrpld/app/ledger/InboundLedger.h>
4#include <xrpld/app/ledger/InboundLedgers.h>
5#include <xrpld/app/ledger/LedgerMaster.h>
6#include <xrpld/app/ledger/LedgerReplayer.h>
7#include <xrpld/app/ledger/detail/TimeoutCounter.h>
8#include <xrpld/app/main/Application.h>
9#include <xrpld/overlay/Peer.h>
10#include <xrpld/overlay/PeerSet.h>
11
12#include <xrpl/basics/Log.h>
13#include <xrpl/basics/base_uint.h>
14#include <xrpl/beast/utility/instrumentation.h>
15#include <xrpl/core/Job.h>
16#include <xrpl/ledger/entries/LedgerHashesEntry.h>
17#include <xrpl/protocol/Indexes.h>
18#include <xrpl/protocol/SField.h>
19#include <xrpl/shamap/SHAMapItem.h>
20
21#include <boost/smart_ptr/intrusive_ptr.hpp>
22
23#include <xrpl.pb.h>
24
25#include <cstddef>
26#include <cstdint>
27#include <memory>
28#include <utility>
29#include <vector>
30
31namespace xrpl {
32
34 Application& app,
35 InboundLedgers& inboundLedgers,
36 UInt256 const& ledgerHash,
39 app,
40 ledgerHash,
41 ledger_replay_parameters::kSubTaskTimeout,
42 {.jobType = JtReplayTask,
43 .jobName = "SkipListAcq",
45 app.getJournal("LedgerReplaySkipList"))
46 , inboundLedgers_(inboundLedgers)
47 , peerSet_(std::move(peerSet))
48{
49 JLOG(journal_.trace()) << "Create " << hash_;
50}
51
53{
54 JLOG(journal_.trace()) << "Destroy " << hash_;
55}
56
57void
59{
61 if (!isDone())
62 {
63 trigger(numPeers, sl);
64 setTimer(sl);
65 }
66}
67
68void
70{
71 if (auto const l = app_.getLedgerMaster().getLedgerByHash(hash_); l)
72 {
73 JLOG(journal_.trace()) << "existing ledger " << hash_;
74 retrieveSkipList(l, sl);
75 return;
76 }
77
78 if (!fallBack_)
79 {
80 peerSet_->addPeers(
81 limit,
82 [this](auto peer) {
83 return peer->supportsFeature(ProtocolFeature::LedgerReplay) &&
84 peer->hasLedger(hash_, 0);
85 },
86 [this](auto peer) {
87 if (peer->supportsFeature(ProtocolFeature::LedgerReplay))
88 {
89 JLOG(journal_.trace()) << "Add a peer " << peer->id() << " for " << hash_;
90 protocol::TMProofPathRequest request;
91 request.set_ledgerhash(hash_.data(), hash_.size());
92 request.set_key(keylet::skip().key.data(), keylet::skip().key.size());
93 request.set_type(protocol::TMLedgerMapType::lmACCOUNT_STATE);
94 peerSet_->sendRequest(request, peer);
95 }
96 else
97 {
98 JLOG(journal_.trace())
99 << "Add a no feature peer " << peer->id() << " for " << hash_;
101 {
102 JLOG(journal_.debug()) << "Fall back for " << hash_;
104 fallBack_ = true;
105 }
106 }
107 });
108 }
109
110 if (fallBack_)
112}
113
114void
116{
117 JLOG(journal_.trace()) << "timeouts_=" << timeouts_ << " for " << hash_;
119 {
120 failed_ = true;
121 JLOG(journal_.debug()) << "too many timeouts " << hash_;
122 notify(sl);
123 }
124 else
125 {
126 trigger(1, sl);
127 }
128}
129
135
136void
138 std::uint32_t ledgerSeq,
139 boost::intrusive_ptr<SHAMapItem const> const& item)
140{
141 XRPL_ASSERT(ledgerSeq != 0 && item, "xrpl::SkipListAcquire::processData : valid inputs");
143 if (isDone())
144 return;
145
146 JLOG(journal_.trace()) << "got data for " << hash_;
147 try
148 {
149 if (auto sle = std::make_shared<SLE>(SerialIter{item->slice()}, item->key()); sle)
150 {
151 if (auto const& skipList = sle->getFieldV256(sfHashes).value(); !skipList.empty())
152 onSkipListAcquired(skipList, ledgerSeq, sl);
153 return;
154 }
155 }
156 catch (...) // NOLINT(bugprone-empty-catch)
157 {
158 }
159
160 failed_ = true;
161 JLOG(journal_.error()) << "failed to retrieve Skip list from verified data " << hash_;
162 notify(sl);
163}
164
165void
167{
169 dataReadyCallbacks_.emplace_back(std::move(cb));
170 if (isDone())
171 {
172 JLOG(journal_.debug()) << "task added to a finished SkipListAcquire " << hash_;
173 notify(sl);
174 }
175}
176
179{
180 ScopedLockType const sl(mtx_);
181 return data_;
182}
183
184void
186{
187 if (LedgerHashesEntryR const hashIndex(*ledger, journal_);
188 hashIndex && hashIndex->isFieldPresent(sfHashes))
189 {
190 auto const& slist = hashIndex->getFieldV256(sfHashes).value();
191 if (!slist.empty())
192 {
193 onSkipListAcquired(slist, ledger->seq(), sl);
194 return;
195 }
196 }
197
198 failed_ = true;
199 JLOG(journal_.error()) << "failed to retrieve Skip list from a ledger " << hash_;
200 notify(sl);
201}
202
203void
205 std::vector<UInt256> const& skipList,
206 std::uint32_t ledgerSeq,
207 ScopedLockType& sl)
208{
209 complete_ = true;
210 data_ = std::make_shared<SkipListData>(ledgerSeq, skipList);
211 JLOG(journal_.debug()) << "Skip list acquired " << hash_;
212 notify(sl);
213}
214
215void
217{
218 XRPL_ASSERT(isDone(), "xrpl::SkipListAcquire::notify : is done");
221 auto const good = !failed_;
222 sl.unlock();
223
224 for (auto& cb : toCall)
225 {
226 cb(good, hash_);
227 }
228
229 sl.lock();
230}
231
232} // namespace xrpl
Manages the lifetime of inbound ledgers.
void addDataCallback(OnSkipListDataCB &&cb)
Add a callback that will be called when the skipList is ready or failed.
std::uint32_t noFeaturePeerCount_
std::weak_ptr< TimeoutCounter > pmDowncast() override
Return a weak pointer to this.
std::unique_ptr< PeerSet > peerSet_
std::shared_ptr< SkipListData const > data_
void trigger(std::size_t limit, ScopedLockType &sl)
Trigger another round.
void notify(ScopedLockType &sl)
Call the OnSkipListDataCB callbacks.
std::function< void(bool successful, UInt256 const &hash)> OnSkipListDataCB
A callback used to notify that the SkipList is ready or failed.
std::vector< OnSkipListDataCB > dataReadyCallbacks_
std::shared_ptr< SkipListData const > getData() const
void onTimer(bool progress, ScopedLockType &peerSetLock) override
Hook called from invokeOnTimer().
InboundLedgers & inboundLedgers_
void onSkipListAcquired(std::vector< UInt256 > const &skipList, std::uint32_t ledgerSeq, ScopedLockType &sl)
Process the skip list.
SkipListAcquire(Application &app, InboundLedgers &inboundLedgers, UInt256 const &ledgerHash, std::unique_ptr< PeerSet > peerSet)
Constructor.
void retrieveSkipList(std::shared_ptr< Ledger const > const &ledger, ScopedLockType &sl)
Retrieve the skip list from the ledger.
void processData(std::uint32_t ledgerSeq, boost::intrusive_ptr< SHAMapItem const > const &item)
Process the data extracted from a peer's reply.
void init(int numPeers)
Start the SkipListAcquire task.
std::recursive_mutex mtx_
std::unique_lock< std::recursive_mutex > ScopedLockType
UInt256 const hash_
The hash of the object (in practice, always a ledger) we are trying to fetch.
TimeoutCounter(Application &app, UInt256 const &targetHash, std::chrono::milliseconds timeoutInterval, QueueJobParameter &&jobParameter, beast::Journal journal)
beast::Journal journal_
void setTimer(ScopedLockType &)
Schedule a call to queueJob() after timerInterval_.
std::chrono::milliseconds timerInterval_
The minimum time to wait between calls to execute().
T lock(T... args)
T make_shared(T... args)
Keylet const & skip() noexcept
The index of the "short" skip list.
Definition Indexes.cpp:232
constexpr std::uint32_t kSubTaskMaxTimeouts
constexpr std::uint32_t kMaxQueuedTasks
Use hash_* containers for keys that do not need a cryptographically secure hashing algorithm.
Definition algorithm.h:5
BaseUInt< 256 > UInt256
Definition base_uint.h:580
@ JtReplayTask
Definition Job.h:47
LedgerHashesEntry< ReadView > LedgerHashesEntryR
T swap(T... args)
T unlock(T... args)