xrpld
Toggle main menu visibility
Loading...
Searching...
No Matches
xrpld
app
ledger
detail
TransactionAcquire.cpp
1
#include <xrpld/app/ledger/detail/TransactionAcquire.h>
2
3
#include <xrpld/app/ledger/ConsensusTransSetSF.h>
4
#include <xrpld/app/ledger/InboundTransactions.h>
5
#include <xrpld/app/ledger/detail/TimeoutCounter.h>
6
#include <xrpld/app/main/Application.h>
7
#include <xrpld/overlay/PeerSet.h>
8
9
#include <xrpl/basics/Log.h>
10
#include <xrpl/basics/base_uint.h>
11
#include <xrpl/core/Job.h>
12
#include <xrpl/server/NetworkOPs.h>
13
#include <xrpl/shamap/SHAMap.h>
14
#include <xrpl/shamap/SHAMapAddNode.h>
15
#include <xrpl/shamap/SHAMapMissingNode.h>
16
#include <xrpl/shamap/SHAMapTreeNode.h>
17
18
#include <xrpl.pb.h>
19
20
#include <
algorithm
>
21
#include <
cstddef
>
22
#include <
exception
>
23
#include <
memory
>
24
#include <
utility
>
25
#include <
vector
>
26
27
namespace
xrpl
{
28
29
using namespace
std::chrono_literals;
30
31
// Timeout interval in milliseconds
32
constexpr
auto
kTxAcquireTimeout
= 250ms;
33
34
static
constexpr
auto
kNormTimeouts
= 4;
35
static
constexpr
auto
kMaxTimeouts
= 20;
36
37
TransactionAcquire::TransactionAcquire
(
38
Application
& app,
39
UInt256
const
& hash,
40
std::unique_ptr<PeerSet>
peerSet)
41
:
TimeoutCounter
(
42
app,
43
hash,
44
kTxAcquireTimeout
,
45
{.jobType =
JtTxnData
, .jobName =
"TxAcq"
, .jobLimit = {}},
46
app.getJournal(
"TransactionAcquire"
))
47
, peerSet_(std::move(peerSet))
48
{
49
map_ =
std::make_shared<SHAMap>
(
SHAMapType::TRANSACTION
, hash, app_.getNodeFamily());
50
map_->setUnbacked();
51
}
52
53
void
54
TransactionAcquire::done
()
55
{
56
// We hold a PeerSet lock and so cannot do real work here
57
58
if
(
failed_
)
59
{
60
JLOG(
journal_
.debug()) <<
"Failed to acquire TX set "
<<
hash_
;
61
}
62
else
63
{
64
JLOG(
journal_
.debug()) <<
"Acquired TX set "
<<
hash_
;
65
map_
->setImmutable();
66
67
UInt256
const
& hash(
hash_
);
68
std::shared_ptr<SHAMap>
const
& map(
map_
);
69
auto
const
pap = &
app_
;
70
// Note that, when we're in the process of shutting down, addJob()
71
// may reject the request. If that happens then giveSet() will
72
// not be called. That's fine. According to David the giveSet() call
73
// just updates the consensus and related structures when we acquire
74
// a transaction set. No need to update them if we're shutting down.
75
app_
.getJobQueue().addJob(
JtTxnData
,
"ComplAcquire"
, [pap, hash, map]() {
76
pap->getInboundTransactions().giveSet(hash, map,
true
);
77
});
78
}
79
}
80
81
void
82
TransactionAcquire::onTimer
(
bool
progress,
ScopedLockType
& psl)
83
{
84
if
(
timeouts_
>
kMaxTimeouts
)
85
{
86
failed_
=
true
;
87
done
();
88
return
;
89
}
90
91
if
(
timeouts_
>=
kNormTimeouts
)
92
trigger
(
nullptr
);
93
94
addPeers
(1);
95
}
96
97
std::weak_ptr<TimeoutCounter>
98
TransactionAcquire::pmDowncast
()
99
{
100
return
shared_from_this
();
101
}
102
103
void
104
TransactionAcquire::trigger
(
std::shared_ptr<Peer>
const
& peer)
105
{
106
if
(
complete_
)
107
{
108
JLOG(
journal_
.info()) <<
"trigger after complete"
;
109
return
;
110
}
111
if
(
failed_
)
112
{
113
JLOG(
journal_
.info()) <<
"trigger after fail"
;
114
return
;
115
}
116
117
if
(!
haveRoot_
)
118
{
119
JLOG(
journal_
.trace()) <<
"TransactionAcquire::trigger "
<< (peer ?
"havePeer"
:
"noPeer"
)
120
<<
" no root"
;
121
protocol::TMGetLedger tmGL;
122
tmGL.set_ledgerhash(
hash_
.begin(),
hash_
.size());
123
tmGL.set_itype(protocol::liTS_CANDIDATE);
124
tmGL.set_querydepth(3);
// We probably need the whole thing
125
126
if
(
timeouts_
!= 0)
127
tmGL.set_querytype(protocol::qtINDIRECT);
128
129
*(tmGL.add_nodeids()) =
SHAMapNodeID
().
getRawString
();
130
peerSet_
->sendRequest(tmGL, peer);
131
}
132
else
if
(!
map_
->isValid())
133
{
134
failed_
=
true
;
135
done
();
136
}
137
else
138
{
139
ConsensusTransSetSF
sf(
app_
,
app_
.getTempNodeCache());
140
auto
nodes =
map_
->getMissingNodes(256, &sf);
141
142
if
(nodes.empty())
143
{
144
if
(
map_
->isValid())
145
{
146
complete_
=
true
;
147
}
148
else
149
{
150
failed_
=
true
;
151
}
152
153
done
();
154
return
;
155
}
156
157
protocol::TMGetLedger tmGL;
158
tmGL.set_ledgerhash(
hash_
.begin(),
hash_
.size());
159
tmGL.set_itype(protocol::liTS_CANDIDATE);
160
161
if
(
timeouts_
!= 0)
162
tmGL.set_querytype(protocol::qtINDIRECT);
163
164
for
(
auto
const
& node : nodes)
165
{
166
*tmGL.add_nodeids() = node.first.getRawString();
167
}
168
peerSet_
->sendRequest(tmGL, peer);
169
}
170
}
171
172
SHAMapAddNode
173
TransactionAcquire::takeNodes
(
174
std::vector
<
std::pair<SHAMapNodeID, SHAMapTreeNodePtr>
> data,
175
std::shared_ptr<Peer>
const
& peer)
176
{
177
ScopedLockType
const
sl(
mtx_
);
178
179
if
(
complete_
)
180
{
181
JLOG(
journal_
.trace()) <<
"TX set complete"
;
182
return
SHAMapAddNode
();
183
}
184
185
if
(
failed_
)
186
{
187
JLOG(
journal_
.trace()) <<
"TX set failed"
;
188
return
SHAMapAddNode
();
189
}
190
191
try
192
{
193
if
(data.empty())
194
return
SHAMapAddNode::invalid
();
195
196
ConsensusTransSetSF
sf(
app_
,
app_
.getTempNodeCache());
197
198
for
(
auto
& d : data)
199
{
200
if
(d.first.isRoot())
201
{
202
if
(
haveRoot_
)
203
{
204
JLOG(
journal_
.debug()) <<
"Got root TXS node, already have it"
;
205
}
206
else
if
(!
map_
->addRootNode(
SHAMapHash
{hash_}, std::move(d.second),
nullptr
)
207
.isGood())
208
{
209
JLOG(
journal_
.warn()) <<
"TX acquire got bad root node for TX set "
<<
hash_
210
<<
" from peer "
<< peer->id();
211
return
SHAMapAddNode::invalid
();
212
}
213
else
214
{
215
haveRoot_
=
true
;
216
}
217
}
218
else
if
(!
map_
->addKnownNode(d.first, std::move(d.second), &sf).isGood())
219
{
220
JLOG(
journal_
.warn()) <<
"TX acquire got bad non-root node "
<< d.first
221
<<
" for TX set "
<<
hash_
<<
" from peer "
<< peer->id();
222
return
SHAMapAddNode::invalid
();
223
}
224
}
225
226
trigger
(peer);
227
progress_
=
true
;
228
return
SHAMapAddNode::useful
();
229
}
230
catch
(
std::exception
const
& ex)
231
{
232
JLOG(
journal_
.error()) <<
"Peer "
<< peer->id()
233
<<
" sent us junky transaction node data: "
<< ex.
what
();
234
return
SHAMapAddNode::invalid
();
235
}
236
}
237
238
void
239
TransactionAcquire::addPeers
(
std::size_t
limit)
240
{
241
peerSet_
->addPeers(
242
limit,
243
[
this
](
auto
peer) {
return
peer->hasTxSet(
hash_
); },
244
[
this
](
auto
peer) {
trigger
(peer); });
245
}
246
247
void
248
TransactionAcquire::init
(
int
numPeers)
249
{
250
ScopedLockType
sl(
mtx_
);
251
252
addPeers
(numPeers);
253
254
setTimer
(sl);
255
}
256
257
void
258
TransactionAcquire::stillNeed
()
259
{
260
ScopedLockType
const
sl(
mtx_
);
261
262
timeouts_
=
std::min<int>
(
timeouts_
,
kNormTimeouts
);
263
failed_
=
false
;
264
}
265
266
}
// namespace xrpl
algorithm
xrpl::Application
Definition
Application.h:94
xrpl::ConsensusTransSetSF
Definition
ConsensusTransSetSF.h:23
xrpl::SHAMapAddNode
Definition
SHAMapAddNode.h:9
xrpl::SHAMapAddNode::useful
static SHAMapAddNode useful()
Definition
SHAMapAddNode.h:124
xrpl::SHAMapAddNode::invalid
static SHAMapAddNode invalid()
Definition
SHAMapAddNode.h:130
xrpl::SHAMapHash
Definition
SHAMapHash.h:16
xrpl::SHAMapNodeID
Identifies a node inside a SHAMap.
Definition
SHAMapNodeID.h:20
xrpl::SHAMapNodeID::getRawString
std::string getRawString() const
Definition
libxrpl/shamap/SHAMapNodeID.cpp:85
xrpl::TimeoutCounter::mtx_
std::recursive_mutex mtx_
Definition
TimeoutCounter.h:121
xrpl::TimeoutCounter::ScopedLockType
std::unique_lock< std::recursive_mutex > ScopedLockType
Definition
TimeoutCounter.h:71
xrpl::TimeoutCounter::hash_
UInt256 const hash_
The hash of the object (in practice, always a ledger) we are trying to fetch.
Definition
TimeoutCounter.h:127
xrpl::TimeoutCounter::TimeoutCounter
TimeoutCounter(Application &app, UInt256 const &targetHash, std::chrono::milliseconds timeoutInterval, QueueJobParameter &&jobParameter, beast::Journal journal)
Definition
TimeoutCounter.cpp:21
xrpl::TimeoutCounter::progress_
bool progress_
Whether forward progress has been made.
Definition
TimeoutCounter.h:134
xrpl::TimeoutCounter::app_
Application & app_
Definition
TimeoutCounter.h:119
xrpl::TimeoutCounter::journal_
beast::Journal journal_
Definition
TimeoutCounter.h:120
xrpl::TimeoutCounter::complete_
bool complete_
Definition
TimeoutCounter.h:129
xrpl::TimeoutCounter::failed_
bool failed_
Definition
TimeoutCounter.h:130
xrpl::TimeoutCounter::setTimer
void setTimer(ScopedLockType &)
Schedule a call to queueJob() after timerInterval_.
Definition
TimeoutCounter.cpp:40
xrpl::TimeoutCounter::timeouts_
int timeouts_
Definition
TimeoutCounter.h:128
xrpl::TransactionAcquire::haveRoot_
bool haveRoot_
Definition
TransactionAcquire.h:47
xrpl::TransactionAcquire::onTimer
void onTimer(bool progress, ScopedLockType &peerSetLock) override
Hook called from invokeOnTimer().
Definition
TransactionAcquire.cpp:82
xrpl::TransactionAcquire::done
void done()
Definition
TransactionAcquire.cpp:54
xrpl::TransactionAcquire::init
void init(int startPeers)
Definition
TransactionAcquire.cpp:248
xrpl::TransactionAcquire::takeNodes
SHAMapAddNode takeNodes(std::vector< std::pair< SHAMapNodeID, SHAMapTreeNodePtr > > data, std::shared_ptr< Peer > const &peer)
Definition
TransactionAcquire.cpp:173
xrpl::TransactionAcquire::map_
std::shared_ptr< SHAMap > map_
Definition
TransactionAcquire.h:46
xrpl::TransactionAcquire::addPeers
void addPeers(std::size_t limit)
Definition
TransactionAcquire.cpp:239
xrpl::TransactionAcquire::TransactionAcquire
TransactionAcquire(Application &app, UInt256 const &hash, std::unique_ptr< PeerSet > peerSet)
Definition
TransactionAcquire.cpp:37
xrpl::TransactionAcquire::trigger
void trigger(std::shared_ptr< Peer > const &)
Definition
TransactionAcquire.cpp:104
xrpl::TransactionAcquire::pmDowncast
std::weak_ptr< TimeoutCounter > pmDowncast() override
Return a weak pointer to this.
Definition
TransactionAcquire.cpp:98
xrpl::TransactionAcquire::stillNeed
void stillNeed()
Definition
TransactionAcquire.cpp:258
xrpl::TransactionAcquire::peerSet_
std::unique_ptr< PeerSet > peerSet_
Definition
TransactionAcquire.h:48
cstddef
exception
std::make_shared
T make_shared(T... args)
memory
std::min
T min(T... args)
xrpl
Use hash_* containers for keys that do not need a cryptographically secure hashing algorithm.
Definition
algorithm.h:5
xrpl::UInt256
BaseUInt< 256 > UInt256
Definition
base_uint.h:580
xrpl::JtTxnData
@ JtTxnData
Definition
Job.h:55
xrpl::kNormTimeouts
static constexpr auto kNormTimeouts
Definition
TransactionAcquire.cpp:34
xrpl::kMaxTimeouts
static constexpr auto kMaxTimeouts
Definition
TransactionAcquire.cpp:35
xrpl::SHAMapType::TRANSACTION
@ TRANSACTION
Definition
SHAMapMissingNode.h:14
xrpl::kTxAcquireTimeout
constexpr auto kTxAcquireTimeout
Definition
TransactionAcquire.cpp:32
std::pair
std::enable_shared_from_this< TransactionAcquire >::shared_from_this
T shared_from_this(T... args)
std::shared_ptr
std::size_t
std::unique_ptr
utility
vector
std::weak_ptr
std::exception::what
T what(T... args)
Generated by
1.17.0