xrpld
Loading...
Searching...
No Matches
InboundLedger.cpp
1#include <xrpld/app/ledger/InboundLedger.h>
2
3#include <xrpld/app/ledger/AccountStateSF.h>
4#include <xrpld/app/ledger/InboundLedgers.h>
5#include <xrpld/app/ledger/LedgerMaster.h>
6#include <xrpld/app/ledger/LedgerNodeHelpers.h>
7#include <xrpld/app/ledger/TransactionStateSF.h>
8#include <xrpld/app/ledger/detail/TimeoutCounter.h>
9#include <xrpld/app/main/Application.h>
10#include <xrpld/overlay/Message.h>
11#include <xrpld/overlay/Overlay.h>
12#include <xrpld/overlay/PeerSet.h>
13
14#include <xrpl/basics/Blob.h>
15#include <xrpl/basics/Log.h>
16#include <xrpl/basics/Slice.h>
17#include <xrpl/basics/base_uint.h>
18#include <xrpl/beast/utility/instrumentation.h>
19#include <xrpl/core/Job.h>
20#include <xrpl/core/JobQueue.h>
21#include <xrpl/json/json_value.h>
22#include <xrpl/ledger/entries/FeeSettingsEntry.h>
23#include <xrpl/nodestore/Database.h>
24#include <xrpl/nodestore/NodeObject.h>
25#include <xrpl/protocol/HashPrefix.h>
26#include <xrpl/protocol/Indexes.h> // IWYU pragma: keep
27#include <xrpl/protocol/LedgerHeader.h>
28#include <xrpl/protocol/Rules.h>
29#include <xrpl/protocol/Serializer.h>
30#include <xrpl/protocol/SystemParameters.h> // IWYU pragma: keep
31#include <xrpl/protocol/jss.h>
32#include <xrpl/resource/Fees.h>
33#include <xrpl/shamap/SHAMapNodeID.h>
34#include <xrpl/shamap/SHAMapSyncFilter.h>
35
36#include <boost/iterator/function_output_iterator.hpp>
37
38#include <xrpl.pb.h>
39
40#include <algorithm>
41#include <chrono>
42#include <cstddef>
43#include <cstdint>
44#include <exception>
45#include <memory>
46#include <mutex>
47#include <random>
48#include <sstream>
49#include <string>
50#include <string_view>
51#include <tuple>
52#include <unordered_map>
53#include <utility>
54#include <vector>
55
56namespace xrpl {
57
58using namespace std::chrono_literals;
59
60static constexpr auto kPeerCountStart = 5; // Number of peers to start with
61static constexpr auto kPeerCountAdd = 3; // Number of peers to add on a timeout
62static constexpr auto kLedgerTimeoutRetriesMax = 6; // how many timeouts before we give up
63static constexpr auto kLedgerBecomeAggressiveThreshold =
64 4; // how many timeouts before we get aggressive
65static constexpr auto kMissingNodesFind = 256; // Number of nodes to find initially
66static constexpr auto kReqNodesReply = 128; // Number of nodes to request for a reply
67static constexpr auto kReqNodes = 12; // Number of nodes to request blindly
68
69// millisecond for each ledger timeout
70constexpr auto kLedgerAcquireTimeout = 3000ms;
71
73 Application& app,
74 UInt256 const& hash,
75 std::uint32_t seq,
76 Reason reason,
77 ClockType& clock,
80 app,
81 hash,
83 {.jobType = JtLedgerData, .jobName = "InboundLedger", .jobLimit = 5},
84 app.getJournal("InboundLedger"))
85 , clock_(clock)
86 , seq_(seq)
87 , reason_(reason)
88 , peerSet_(std::move(peerSet))
89{
90 JLOG(journal_.trace()) << "Acquiring ledger " << hash_;
91 touch();
92}
93
94void
96{
98 collectionLock.unlock();
99
100 tryDB(app_.getNodeFamily().db());
101 if (failed_)
102 return;
103
104 if (!complete_)
105 {
106 addPeers();
107 queueJob(sl);
108 return;
109 }
110
111 JLOG(journal_.debug()) << "Acquiring ledger we already have in "
112 << " local store. " << hash_;
113 XRPL_ASSERT(
115 "xrpl::InboundLedger::init : valid ledger fees");
116 ledger_->setImmutable();
117
119 return;
120
121 app_.getLedgerMaster().storeLedger(ledger_);
122
123 // Check if this could be a newer fully-validated ledger
125 app_.getLedgerMaster().checkAccept(ledger_);
126}
127
130{
131 auto const& peerIds = peerSet_->getPeerIds();
133 peerIds, [this](auto id) { return (app_.getOverlay().findPeerByShortID(id) != nullptr); });
134}
135
136void
138{
139 ScopedLockType const sl(mtx_);
140
141 // If we didn't know the sequence number, but now do, save it
142 if ((seq != 0) && (seq_ == 0))
143 seq_ = seq;
144
145 // Prevent this from being swept
146 touch();
147}
148
149bool
151{
152 ScopedLockType const sl(mtx_);
153 if (!isDone())
154 {
155 if (ledger_)
156 {
157 tryDB(ledger_->stateMap().family().db());
158 }
159 else
160 {
161 tryDB(app_.getNodeFamily().db());
162 }
163 if (failed_ || complete_)
164 {
165 done();
166 return true;
167 }
168 }
169 return false;
170}
171
173{
174 // Save any received AS data not processed. It could be useful
175 // for populating a different ledger
176 for (auto& entry : receivedData_)
177 {
178 if (entry.second->type() == protocol::liAS_NODE)
179 app_.getInboundLedgers().gotStaleData(entry.second);
180 }
181 if (!isDone())
182 {
183 JLOG(journal_.debug()) << "Acquire " << hash_ << " abort "
184 << ((timeouts_ == 0) ? std::string()
185 : (std::string("timeouts:") +
187 << stats_.get();
188 }
189}
190
192neededHashes(UInt256 const& root, SHAMap& map, int max, SHAMapSyncFilter const* filter)
193{
195
196 if (!root.isZero())
197 {
198 if (map.getHash().isZero())
199 {
200 ret.push_back(root);
201 }
202 else
203 {
204 auto mn = map.getMissingNodes(max, filter);
205 ret.reserve(mn.size());
206 for (auto const& n : mn)
207 ret.push_back(n.second);
208 }
209 }
210
211 return ret;
212}
213
216{
217 return neededHashes(ledger_->header().txHash, ledger_->txMap(), max, filter);
218}
219
222{
223 return neededHashes(ledger_->header().accountHash, ledger_->stateMap(), max, filter);
224}
225
226// See how much of the ledger data is stored locally
227// Data found in a fetch pack will be stored
228void
230{
231 if (!haveHeader_)
232 {
233 auto makeLedger = [&, this](Blob const& data) {
234 JLOG(journal_.trace()) << "Ledger header found in fetch pack";
235 Rules const rules{app_.config().features};
237 deserializePrefixedHeader(makeSlice(data)), rules, app_.getNodeFamily());
238 if (ledger_->header().hash != hash_ || (seq_ != 0 && seq_ != ledger_->header().seq))
239 {
240 // We know for a fact the ledger can never be acquired
241 JLOG(journal_.warn())
242 << "hash " << hash_ << " seq " << std::to_string(seq_) << " cannot be a ledger";
243 ledger_.reset();
244 failed_ = true;
245 }
246 };
247
248 // Try to fetch the ledger header from the DB
249 if (auto nodeObject = srcDB.fetchNodeObject(hash_, seq_))
250 {
251 JLOG(journal_.trace()) << "Ledger header found in local store";
252
253 makeLedger(nodeObject->getData());
254 if (failed_)
255 return;
256
257 // Store the ledger header if the source and destination differ
258 auto& dstDB{ledger_->stateMap().family().db()};
259 if (std::addressof(dstDB) != std::addressof(srcDB))
260 {
261 Blob blob{nodeObject->getData()};
262 dstDB.store(NodeObjectType::Ledger, std::move(blob), hash_, ledger_->header().seq);
263 }
264 }
265 else
266 {
267 // Try to fetch the ledger header from a fetch pack
268 auto data = app_.getLedgerMaster().getFetchPack(hash_);
269 if (!data)
270 return;
271
272 JLOG(journal_.trace()) << "Ledger header found in fetch pack";
273
274 makeLedger(*data);
275 if (failed_)
276 return;
277
278 // Store the ledger header in the ledger's database
279 ledger_->stateMap().family().db().store(
280 NodeObjectType::Ledger, std::move(*data), hash_, ledger_->header().seq);
281 }
282
283 if (seq_ == 0)
284 seq_ = ledger_->header().seq;
285 ledger_->stateMap().setLedgerSeq(seq_);
286 ledger_->txMap().setLedgerSeq(seq_);
287 haveHeader_ = true;
288 }
289
291 {
292 if (ledger_->header().txHash.isZero())
293 {
294 JLOG(journal_.trace()) << "No TXNs to fetch";
295 haveTransactions_ = true;
296 }
297 else
298 {
299 TransactionStateSF filter(ledger_->txMap().family().db(), app_.getLedgerMaster());
300 if (ledger_->txMap().fetchRoot(SHAMapHash{ledger_->header().txHash}, &filter))
301 {
302 if (neededTxHashes(1, &filter).empty())
303 {
304 JLOG(journal_.trace()) << "Had full txn map locally";
305 haveTransactions_ = true;
306 }
307 }
308 }
309 }
310
311 if (!haveState_)
312 {
313 if (ledger_->header().accountHash.isZero())
314 {
315 JLOG(journal_.fatal()) << "We are acquiring a ledger with a zero account hash";
316 failed_ = true;
317 return;
318 }
319 AccountStateSF filter(ledger_->stateMap().family().db(), app_.getLedgerMaster());
320 if (ledger_->stateMap().fetchRoot(SHAMapHash{ledger_->header().accountHash}, &filter))
321 {
322 if (neededStateHashes(1, &filter).empty())
323 {
324 JLOG(journal_.trace()) << "Had full AS map locally";
325 haveState_ = true;
326 }
327 }
328 }
329
331 {
332 JLOG(journal_.debug()) << "Had everything locally";
333 complete_ = true;
334 XRPL_ASSERT(
336 "xrpl::InboundLedger::tryDB : valid ledger fees");
337 ledger_->setImmutable();
338 }
339}
340
344void
346{
347 recentNodes_.clear();
348
349 if (isDone())
350 {
351 JLOG(journal_.info()) << "Already done " << hash_;
352 return;
353 }
354
356 {
357 if (seq_ != 0)
358 {
359 JLOG(journal_.warn()) << timeouts_ << " timeouts for ledger " << seq_;
360 }
361 else
362 {
363 JLOG(journal_.warn()) << timeouts_ << " timeouts for ledger " << hash_;
364 }
365 failed_ = true;
366 done();
367 return;
368 }
369
370 if (!wasProgress)
371 {
372 checkLocal();
373
374 byHash_ = true;
375
376 std::size_t const pc = getPeerCount();
377 JLOG(journal_.debug()) << "No progress(" << pc << ") for ledger " << hash_;
378
379 // addPeers triggers if the reason is not HISTORY
380 // So if the reason IS HISTORY, need to trigger after we add
381 // otherwise, we need to trigger before we add
382 // so each peer gets triggered once
385 addPeers();
388 }
389}
390
394void
396{
397 peerSet_->addPeers(
399 [this](auto peer) { return peer->hasLedger(hash_, seq_); },
400 [this](auto peer) {
401 // For historical nodes, do not trigger too soon
402 // since a fetch pack is probably coming
405 });
406}
407
413
414void
416{
417 if (signaled_)
418 return;
419
420 signaled_ = true;
421 touch();
422
423 JLOG(journal_.debug()) << "Acquire " << hash_ << (failed_ ? " fail " : " ")
424 << ((timeouts_ == 0)
425 ? std::string()
426 : (std::string("timeouts:") + std::to_string(timeouts_) + " "))
427 << stats_.get();
428
429 XRPL_ASSERT(complete_ || failed_, "xrpl::InboundLedger::done : complete or failed");
430
431 if (complete_ && !failed_ && ledger_)
432 {
433 XRPL_ASSERT(
435 "xrpl::InboundLedger::done : valid ledger fees");
436 ledger_->setImmutable();
437 switch (reason_)
438 {
439 case Reason::HISTORY:
440 app_.getInboundLedgers().onLedgerFetched();
441 break;
442 default:
443 app_.getLedgerMaster().storeLedger(ledger_);
444 break;
445 }
446 }
447
448 // We hold the PeerSet lock, so must dispatch
449 app_.getJobQueue().addJob(JtLedgerData, "AcqDone", [self = shared_from_this()]() {
450 if (self->complete_ && !self->failed_)
451 {
452 self->app_.getLedgerMaster().checkAccept(self->getLedger());
453 self->app_.getLedgerMaster().tryAdvance();
454 }
455 else
456 {
457 self->app_.getInboundLedgers().logFailure(self->hash_, self->seq_);
458 }
459 });
460}
461
465void
467{
469
470 if (isDone())
471 {
472 JLOG(journal_.debug()) << "Trigger on ledger: " << hash_ << (complete_ ? " completed" : "")
473 << (failed_ ? " failed" : "");
474 return;
475 }
476
477 if (auto stream = journal_.debug())
478 {
480 ss << "Trigger acquiring ledger " << hash_;
481 if (peer)
482 ss << " from " << peer;
483
484 if (complete_ || failed_)
485 {
486 ss << " complete=" << complete_ << " failed=" << failed_;
487 }
488 else
489 {
490 ss << " header=" << haveHeader_ << " tx=" << haveTransactions_ << " as=" << haveState_;
491 }
492 stream << ss.str();
493 }
494
495 if (!haveHeader_)
496 {
497 tryDB(app_.getNodeFamily().db());
498 if (failed_)
499 {
500 JLOG(journal_.warn()) << " failed local for " << hash_;
501 return;
502 }
503 }
504
505 protocol::TMGetLedger tmGL;
506 tmGL.set_ledgerhash(hash_.begin(), hash_.size());
507
508 if (timeouts_ != 0)
509 {
510 // Be more aggressive if we've timed out at least once
511 tmGL.set_querytype(protocol::qtINDIRECT);
512
514 {
515 auto need = getNeededHashes();
516
517 if (!need.empty())
518 {
519 protocol::TMGetObjectByHash tmBH;
520 bool typeSet = false;
521 tmBH.set_query(true);
522 tmBH.set_ledgerhash(hash_.begin(), hash_.size());
523 for (auto const& p : need)
524 {
525 JLOG(journal_.debug()) << "Want: " << p.second;
526
527 if (!typeSet)
528 {
529 tmBH.set_type(p.first);
530 typeSet = true;
531 }
532
533 if (p.first == tmBH.type())
534 {
535 protocol::TMIndexedObject* io = tmBH.add_objects();
536 io->set_hash(p.second.begin(), p.second.size());
537 if (seq_ != 0)
538 io->set_ledgerseq(seq_);
539 }
540 }
541
542 auto packet = std::make_shared<Message>(tmBH, protocol::mtGET_OBJECTS);
543 auto const& peerIds = peerSet_->getPeerIds();
544 std::ranges::for_each(peerIds, [this, &packet](auto id) {
545 if (auto p = app_.getOverlay().findPeerByShortID(id))
546 {
547 byHash_ = false;
548 p->send(packet);
549 }
550 });
551 }
552 else
553 {
554 JLOG(journal_.info()) << "getNeededHashes says acquire is complete";
555 haveHeader_ = true;
556 haveTransactions_ = true;
557 haveState_ = true;
558 complete_ = true;
559 }
560 }
561 }
562
563 // We can't do much without the header data because we don't know the
564 // state or transaction root hashes.
565 if (!haveHeader_ && !failed_)
566 {
567 tmGL.set_itype(protocol::liBASE);
568 if (seq_ != 0)
569 tmGL.set_ledgerseq(seq_);
570 JLOG(journal_.trace()) << "Sending header request to "
571 << (peer ? "selected peer" : "all peers");
572 peerSet_->sendRequest(tmGL, peer);
573 return;
574 }
575
576 if (ledger_)
577 tmGL.set_ledgerseq(ledger_->header().seq);
578
579 if (reason != TriggerReason::Reply)
580 {
581 // If we're querying blind, don't query deep
582 tmGL.set_querydepth(0);
583 }
584 else if (peer && peer->isHighLatency())
585 {
586 // If the peer has high latency, query extra deep
587 tmGL.set_querydepth(2);
588 }
589 else
590 {
591 tmGL.set_querydepth(1);
592 }
593
594 // Get the state data first because it's the most likely to be useful
595 // if we wind up abandoning this fetch.
596 if (haveHeader_ && !haveState_ && !failed_)
597 {
598 XRPL_ASSERT(
599 ledger_,
600 "xrpl::InboundLedger::trigger : non-null ledger to read state "
601 "from");
602
603 if (!ledger_->stateMap().isValid())
604 {
605 failed_ = true;
606 }
607 else if (ledger_->stateMap().getHash().isZero())
608 {
609 // we need the root node
610 tmGL.set_itype(protocol::liAS_NODE);
611 *tmGL.add_nodeids() = SHAMapNodeID().getRawString();
612 JLOG(journal_.trace())
613 << "Sending AS root request to " << (peer ? "selected peer" : "all peers");
614 peerSet_->sendRequest(tmGL, peer);
615 return;
616 }
617 else
618 {
619 AccountStateSF filter(ledger_->stateMap().family().db(), app_.getLedgerMaster());
620
621 // Release the lock while we process the large state map
622 sl.unlock();
623 auto nodes = ledger_->stateMap().getMissingNodes(kMissingNodesFind, &filter);
624 sl.lock();
625
626 // Make sure nothing happened while we released the lock
627 if (!failed_ && !complete_ && !haveState_)
628 {
629 if (nodes.empty())
630 {
631 if (!ledger_->stateMap().isValid())
632 {
633 failed_ = true;
634 }
635 else
636 {
637 haveState_ = true;
638
640 complete_ = true;
641 }
642 }
643 else
644 {
645 filterNodes(nodes, reason);
646
647 if (!nodes.empty())
648 {
649 tmGL.set_itype(protocol::liAS_NODE);
650 for (auto const& id : nodes)
651 {
652 *(tmGL.add_nodeids()) = id.first.getRawString();
653 }
654
655 JLOG(journal_.trace()) << "Sending AS node request (" << nodes.size()
656 << ") to " << (peer ? "selected peer" : "all peers");
657 peerSet_->sendRequest(tmGL, peer);
658 return;
659 }
660
661 JLOG(journal_.trace()) << "All AS nodes filtered";
662 }
663 }
664 }
665 }
666
668 {
669 XRPL_ASSERT(
670 ledger_,
671 "xrpl::InboundLedger::trigger : non-null ledger to read "
672 "transactions from");
673
674 if (!ledger_->txMap().isValid())
675 {
676 failed_ = true;
677 }
678 else if (ledger_->txMap().getHash().isZero())
679 {
680 // we need the root node
681 tmGL.set_itype(protocol::liTX_NODE);
682 *(tmGL.add_nodeids()) = SHAMapNodeID().getRawString();
683 JLOG(journal_.trace())
684 << "Sending TX root request to " << (peer ? "selected peer" : "all peers");
685 peerSet_->sendRequest(tmGL, peer);
686 return;
687 }
688 else
689 {
690 TransactionStateSF filter(ledger_->txMap().family().db(), app_.getLedgerMaster());
691
692 auto nodes = ledger_->txMap().getMissingNodes(kMissingNodesFind, &filter);
693
694 if (nodes.empty())
695 {
696 if (!ledger_->txMap().isValid())
697 {
698 failed_ = true;
699 }
700 else
701 {
702 haveTransactions_ = true;
703
704 if (haveState_)
705 complete_ = true;
706 }
707 }
708 else
709 {
710 filterNodes(nodes, reason);
711
712 if (!nodes.empty())
713 {
714 tmGL.set_itype(protocol::liTX_NODE);
715 for (auto const& n : nodes)
716 {
717 *(tmGL.add_nodeids()) = n.first.getRawString();
718 }
719 JLOG(journal_.trace()) << "Sending TX node request (" << nodes.size() << ") to "
720 << (peer ? "selected peer" : "all peers");
721 peerSet_->sendRequest(tmGL, peer);
722 return;
723 }
724
725 JLOG(journal_.trace()) << "All TX nodes filtered";
726 }
727 }
728 }
729
730 if (complete_ || failed_)
731 {
732 JLOG(journal_.debug()) << "Done:" << (complete_ ? " complete" : "")
733 << (failed_ ? " failed " : " ") << ledger_->header().seq;
734 sl.unlock();
735 done();
736 }
737}
738
739void
742 TriggerReason reason)
743{
744 // Sort nodes so that the ones we haven't recently
745 // requested come before the ones we have.
747 nodes, [this](auto const& item) { return recentNodes_.count(item.second) == 0; });
748
749 // If everything is a duplicate we don't want to send
750 // any query at all except on a timeout where we need
751 // to query everyone:
752 if (dup.begin() == nodes.begin())
753 {
754 JLOG(journal_.trace()) << "filterNodes: all duplicates";
755
756 if (reason != TriggerReason::Timeout)
757 {
758 nodes.clear();
759 return;
760 }
761 }
762 else
763 {
764 JLOG(journal_.trace()) << "filterNodes: pruning duplicates";
765
766 nodes.erase(dup.begin(), dup.end());
767 }
768
769 std::size_t const limit = (reason == TriggerReason::Reply) ? kReqNodesReply : kReqNodes;
770
771 if (nodes.size() > limit)
772 nodes.resize(limit);
773
774 for (auto const& n : nodes)
775 recentNodes_.insert(n.second);
776}
777
782// data must not have hash prefix
783bool
785{
786 // Return value: true=normal, false=bad data
787 JLOG(journal_.trace()) << "got header acquiring ledger " << hash_;
788
789 if (complete_ || failed_ || haveHeader_)
790 return true;
791
792 auto* f = &app_.getNodeFamily();
793 Rules const rules{app_.config().features};
795 if (ledger_->header().hash != hash_ || (seq_ != 0 && seq_ != ledger_->header().seq))
796 {
797 JLOG(journal_.warn()) << "Acquire hash mismatch: " << ledger_->header().hash
798 << "!=" << hash_;
799 ledger_.reset();
800 return false;
801 }
802 if (seq_ == 0)
803 seq_ = ledger_->header().seq;
804 ledger_->stateMap().setLedgerSeq(seq_);
805 ledger_->txMap().setLedgerSeq(seq_);
806 haveHeader_ = true;
807
808 Serializer s(data.size() + 4);
810 s.addRaw(data.data(), data.size());
811 f->db().store(NodeObjectType::Ledger, std::move(s.modData()), hash_, seq_);
812
813 if (ledger_->header().txHash.isZero())
814 haveTransactions_ = true;
815
816 if (ledger_->header().accountHash.isZero())
817 haveState_ = true;
818
819 ledger_->txMap().setSynching();
820 ledger_->stateMap().setSynching();
821
822 return true;
823}
824
829void
831 std::shared_ptr<Peer> const& peer,
832 protocol::TMLedgerData const& packet,
833 SHAMapAddNode& san)
834{
835 if (!haveHeader_)
836 {
837 JLOG(journal_.warn()) << "Missing ledger header";
838 san.incInvalid();
839 return;
840 }
841 if (packet.type() == protocol::liTX_NODE)
842 {
844 {
845 san.incDuplicate();
846 return;
847 }
848 }
849 else if (haveState_ || failed_)
850 {
851 san.incDuplicate();
852 return;
853 }
854
855 auto [map, rootHash, filter] =
857 if (packet.type() == protocol::liTX_NODE)
858 {
859 return {
860 ledger_->txMap(),
861 SHAMapHash{ledger_->header().txHash},
863 ledger_->txMap().family().db(), app_.getLedgerMaster())};
864 }
865 return {
866 ledger_->stateMap(),
867 SHAMapHash{ledger_->header().accountHash},
869 ledger_->stateMap().family().db(), app_.getLedgerMaster())};
870 }();
871
872 try
873 {
874 auto const f = filter.get();
875
876 for (auto const& ledgerNode : packet.nodes())
877 {
878 auto treeNode = getTreeNode(ledgerNode.nodedata());
879 if (!treeNode)
880 {
881 JLOG(journal_.warn())
882 << "Got invalid node data for ledger " << hash_ << " from peer " << peer->id();
883 peer->charge(resource::kFeeInvalidData, "ledger_node.node_data invalid");
884 san.incInvalid();
885 return;
886 }
887
888 auto const nodeID = getSHAMapNodeID(ledgerNode, *treeNode);
889 if (!nodeID)
890 {
891 JLOG(journal_.warn())
892 << "Got invalid node id for ledger " << hash_ << " from peer " << peer->id();
893 peer->charge(resource::kFeeInvalidData, "ledger_node.node_id invalid");
894 san.incInvalid();
895 return;
896 }
897
898 auto const result = nodeID->isRoot()
899 ? map.addRootNode(rootHash, std::move(treeNode), f)
900 : map.addKnownNode(*nodeID, std::move(treeNode), f);
901 san += result;
902
903 if (result.isInvalid())
904 {
905 JLOG(journal_.warn()) << "Got invalid node " << *nodeID << " for ledger " << hash_
906 << " from peer " << peer->id();
907 peer->charge(resource::kFeeInvalidData, "ledger_node invalid");
908 return;
909 }
910 }
911 }
912 catch (std::exception const& e)
913 {
914 // If we get here it is not necessarily because the node was bad, so don't charge the peer.
915 JLOG(journal_.error()) << "Could not process node for ledger " << hash_ << " from peer "
916 << peer->id() << ": " << e.what();
917 san.incInvalid();
918 return;
919 }
920
921 if (!map.isSynching())
922 {
923 if (packet.type() == protocol::liTX_NODE)
924 {
925 haveTransactions_ = true;
926 }
927 else
928 {
929 haveState_ = true;
930 }
931
933 {
934 complete_ = true;
935 done();
936 }
937 }
938}
939
944bool
946{
947 if (failed_ || haveState_)
948 {
949 san.incDuplicate();
950 return true;
951 }
952
953 if (!haveHeader_)
954 {
955 // LCOV_EXCL_START
956 UNREACHABLE("xrpl::InboundLedger::takeAsRootNode : no ledger header");
957 return false;
958 // LCOV_EXCL_STOP
959 }
960
961 auto treeNode = getTreeNode(data);
962 if (!treeNode)
963 {
964 JLOG(journal_.warn()) << "Got invalid AS root node data for ledger " << hash_;
965 san.incInvalid();
966 return false;
967 }
968
969 AccountStateSF filter(ledger_->stateMap().family().db(), app_.getLedgerMaster());
970 auto const result = ledger_->stateMap().addRootNode(
971 SHAMapHash{ledger_->header().accountHash}, std::move(treeNode), &filter);
972 san += result;
973 return !result.isInvalid();
974}
975
980bool
982{
984 {
985 san.incDuplicate();
986 return true;
987 }
988
989 if (!haveHeader_)
990 {
991 // LCOV_EXCL_START
992 UNREACHABLE("xrpl::InboundLedger::takeTxRootNode : no ledger header");
993 return false;
994 // LCOV_EXCL_STOP
995 }
996
997 auto treeNode = getTreeNode(data);
998 if (!treeNode)
999 {
1000 JLOG(journal_.warn()) << "Got invalid TX root node data for ledger " << hash_;
1001 san.incInvalid();
1002 return false;
1003 }
1004
1005 TransactionStateSF filter(ledger_->txMap().family().db(), app_.getLedgerMaster());
1006 auto const result = ledger_->txMap().addRootNode(
1007 SHAMapHash{ledger_->header().txHash}, std::move(treeNode), &filter);
1008 san += result;
1009 return !result.isInvalid();
1010}
1011
1014{
1016
1017 if (!haveHeader_)
1018 {
1019 ret.emplace_back(protocol::TMGetObjectByHash::otLEDGER, hash_);
1020 return ret;
1021 }
1022
1023 if (!haveState_)
1024 {
1025 AccountStateSF filter(ledger_->stateMap().family().db(), app_.getLedgerMaster());
1026 for (auto const& h : neededStateHashes(4, &filter))
1027 {
1028 ret.emplace_back(protocol::TMGetObjectByHash::otSTATE_NODE, h);
1029 }
1030 }
1031
1032 if (!haveTransactions_)
1033 {
1034 TransactionStateSF filter(ledger_->txMap().family().db(), app_.getLedgerMaster());
1035 for (auto const& h : neededTxHashes(4, &filter))
1036 {
1037 ret.emplace_back(protocol::TMGetObjectByHash::otTRANSACTION_NODE, h);
1038 }
1039 }
1040
1041 return ret;
1042}
1043
1048bool
1052{
1054
1055 if (isDone())
1056 return false;
1057
1058 receivedData_.emplace_back(peer, data);
1059
1061 return false;
1062
1063 receiveDispatched_ = true;
1064 return true;
1065}
1066
1071// VFALCO NOTE, it is not necessary to pass the entire Peer,
1072// we can get away with just a resource::Consumer endpoint.
1073//
1074// TODO Change peer to Consumer
1075//
1076int
1077InboundLedger::processData(std::shared_ptr<Peer> peer, protocol::TMLedgerData const& packet)
1078{
1079 if (packet.type() == protocol::liBASE)
1080 {
1081 if (packet.nodes().empty())
1082 {
1083 JLOG(journal_.warn()) << peer->id() << ": empty header data";
1084 peer->charge(resource::kFeeMalformedRequest, "ledger_data empty header");
1085 return -1;
1086 }
1087
1088 SHAMapAddNode san;
1089
1090 ScopedLockType const sl(mtx_);
1091
1092 try
1093 {
1094 if (!haveHeader_)
1095 {
1096 if (!takeHeader(packet.nodes(0).nodedata()))
1097 {
1098 JLOG(journal_.warn()) << "Got invalid header data";
1099 peer->charge(resource::kFeeMalformedRequest, "ledger_data invalid header");
1100 return -1;
1101 }
1102
1103 san.incUseful();
1104 }
1105
1106 if (!haveState_ && (packet.nodes().size() > 1) &&
1107 !takeAsRootNode(packet.nodes(1).nodedata(), san))
1108 {
1109 JLOG(journal_.warn()) << "Included AS root invalid for ledger " << hash_
1110 << " from peer " << peer->id();
1111 if (san.isInvalid())
1112 {
1113 peer->charge(resource::kFeeInvalidData, "ledger_data invalid AS root");
1114 return -1;
1115 }
1116 }
1117
1118 if (!haveTransactions_ && (packet.nodes().size() > 2) &&
1119 !takeTxRootNode(packet.nodes(2).nodedata(), san))
1120 {
1121 JLOG(journal_.warn()) << "Included TX root invalid for ledger " << hash_
1122 << " from peer " << peer->id();
1123 if (san.isInvalid())
1124 {
1125 peer->charge(resource::kFeeInvalidData, "ledger_data invalid TX root");
1126 return -1;
1127 }
1128 }
1129 }
1130 catch (std::exception const& ex)
1131 {
1132 JLOG(journal_.warn()) << "Included AS/TX root invalid for ledger " << hash_
1133 << " from peer " << peer->id() << ": " << ex.what();
1134 using namespace std::string_literals;
1135 peer->charge(resource::kFeeInvalidData, "ledger_data "s + ex.what());
1136 return -1;
1137 }
1138
1139 if (san.isUseful())
1140 progress_ = true;
1141
1142 stats_ += san;
1143 return san.getGood();
1144 }
1145
1146 if ((packet.type() == protocol::liTX_NODE) || (packet.type() == protocol::liAS_NODE))
1147 {
1148 if (packet.nodes().empty())
1149 {
1150 JLOG(journal_.info()) << peer->id() << ": response with no nodes";
1151 peer->charge(resource::kFeeMalformedRequest, "ledger_data no nodes");
1152 return -1;
1153 }
1154
1155 ScopedLockType const sl(mtx_);
1156
1157 SHAMapAddNode san;
1158 receiveNode(peer, packet, san);
1159
1160 JLOG(journal_.debug()) << "Ledger "
1161 << ((packet.type() == protocol::liTX_NODE) ? "TX" : "AS")
1162 << " node stats: " << san.get();
1163
1164 // `san` accumulates across the whole packet, so `isInvalid()` (bad_ > 0) does not mean the
1165 // packet had no useful nodes: credit whatever good/useful nodes were sent rather than
1166 // discarding everything because one node in an otherwise-good packet was bad.
1167 // Note: Peer charges for invalid/malformed data are issued from within receiveNode at the
1168 // exact failure site, so the peer is only charged for problems they are responsible for.
1169 if (san.isUseful())
1170 progress_ = true;
1171
1172 stats_ += san;
1173 return san.getGood();
1174 }
1175
1176 return -1;
1177}
1178
1179namespace detail {
1180// Track the amount of useful data that each peer returns
1182{
1183 // Map from peer to amount of useful the peer returned
1185 // The largest amount of useful data that any peer returned
1186 int maxCount = 0;
1187
1188 // Update the data count for a peer
1189 void
1190 update(std::shared_ptr<Peer>&& peer, int dataCount)
1191 {
1192 if (dataCount <= 0)
1193 return;
1194 maxCount = std::max(maxCount, dataCount);
1195 auto i = counts.find(peer);
1196 if (i == counts.end())
1197 {
1198 counts.emplace(std::move(peer), dataCount);
1199 return;
1200 }
1201 i->second = std::max(i->second, dataCount);
1202 }
1203
1204 // Prune all the peers that didn't return enough data.
1205 void
1207 {
1208 // Remove all the peers that didn't return at least half as much data as
1209 // the best peer
1210 auto const thresh = maxCount / 2;
1211 auto i = counts.begin();
1212 while (i != counts.end())
1213 {
1214 if (i->second < thresh)
1215 {
1216 i = counts.erase(i);
1217 }
1218 else
1219 {
1220 ++i;
1221 }
1222 }
1223 }
1224
1225 // call F with the `peer` parameter with a random sample of at most n values
1226 // of the counts vector.
1227 template <class F>
1228 void
1230 {
1231 if (counts.empty())
1232 return;
1233
1234 auto outFunc = [&f](auto&& v) { f(v.first); };
1236#if _MSC_VER
1238 s.reserve(n);
1239 std::sample(counts.begin(), counts.end(), std::back_inserter(s), n, rng);
1240 for (auto& v : s)
1241 {
1242 outFunc(v);
1243 }
1244#else
1246 counts.begin(), counts.end(), boost::make_function_output_iterator(outFunc), n, rng);
1247#endif
1248 }
1249};
1250} // namespace detail
1251
1256void
1258{
1259 // Maximum number of peers to request data from
1260 static constexpr std::size_t kMaxUsefulPeers = 6;
1261
1262 decltype(receivedData_) data;
1263
1264 // Reserve some memory so the first couple iterations don't reallocate
1265 data.reserve(8);
1266
1267 detail::PeerDataCounts dataCounts;
1268
1269 for (;;)
1270 {
1271 data.clear();
1272
1273 {
1275
1276 if (receivedData_.empty())
1277 {
1278 receiveDispatched_ = false;
1279 break;
1280 }
1281
1282 data.swap(receivedData_);
1283 }
1284
1285 for (auto& entry : data)
1286 {
1287 if (auto peer = entry.first.lock())
1288 {
1289 int const count = processData(peer, *(entry.second));
1290 dataCounts.update(std::move(peer), count);
1291 }
1292 }
1293 }
1294
1295 // Select a random sample of the peers that gives us the most nodes that are
1296 // useful
1297 dataCounts.prune();
1298 dataCounts.sampleN(kMaxUsefulPeers, [&](std::shared_ptr<Peer> const& peer) {
1300 });
1301}
1302
1305{
1307
1308 ScopedLockType const sl(mtx_);
1309
1310 ret[jss::hash] = to_string(hash_);
1311
1312 if (complete_)
1313 ret[jss::complete] = true;
1314
1315 if (failed_)
1316 ret[jss::failed] = true;
1317
1318 if (!complete_ && !failed_)
1319 ret[jss::peers] = static_cast<int>(peerSet_->getPeerIds().size());
1320
1321 ret[jss::have_header] = haveHeader_;
1322
1323 if (haveHeader_)
1324 {
1325 ret[jss::have_state] = haveState_;
1326 ret[jss::have_transactions] = haveTransactions_;
1327 }
1328
1329 ret[jss::timeouts] = timeouts_;
1330
1331 if (haveHeader_ && !haveState_)
1332 {
1334 for (auto const& h : neededStateHashes(16, nullptr))
1335 {
1336 hv.append(to_string(h));
1337 }
1338 ret[jss::needed_state_hashes] = hv;
1339 }
1340
1342 {
1344 for (auto const& h : neededTxHashes(16, nullptr))
1345 {
1346 hv.append(to_string(h));
1347 }
1348 ret[jss::needed_transaction_hashes] = hv;
1349 }
1350
1351 return ret;
1352}
1353
1354} // namespace xrpl
T addressof(T... args)
T back_inserter(T... args)
Represents a JSON value.
Definition json_value.h:117
Value & append(Value const &value)
Append value to array at the end.
ValueType type() const
json::Value getJson(int)
Return a json::ValueType::Object.
void tryDB(node_store::Database &srcDB)
void trigger(std::shared_ptr< Peer > const &, TriggerReason)
Request more nodes, perhaps from a specific peer.
std::weak_ptr< TimeoutCounter > pmDowncast() override
Return a weak pointer to this.
void runData()
Process pending TMLedgerData Query the a random sample of the 'best' peers.
void filterNodes(std::vector< std::pair< SHAMapNodeID, UInt256 > > &nodes, TriggerReason reason)
std::size_t getPeerCount() const
void onTimer(bool progress, ScopedLockType &peerSetLock) override
Called with a lock by the PeerSet when the timer expires.
std::set< UInt256 > recentNodes_
void receiveNode(std::shared_ptr< Peer > const &peer, protocol::TMLedgerData const &packet, SHAMapAddNode &san)
Process node data received from a peer Call with a lock.
std::vector< UInt256 > neededStateHashes(int max, SHAMapSyncFilter const *filter) const
SHAMapAddNode stats_
int processData(std::shared_ptr< Peer > peer, protocol::TMLedgerData const &data)
Process one TMLedgerData Returns the number of useful nodes.
bool takeHeader(std::string_view data)
Take ledger header data Call with a lock.
bool takeAsRootNode(std::string_view data, SHAMapAddNode &san)
Process AS root node received from a peer Call with a lock.
std::shared_ptr< Ledger > ledger_
InboundLedger(Application &app, UInt256 const &hash, std::uint32_t seq, Reason reason, ClockType &, std::unique_ptr< PeerSet > peerSet)
std::vector< std::pair< std::weak_ptr< Peer >, std::shared_ptr< protocol::TMLedgerData > > > receivedData_
std::mutex receivedDataLock_
std::unique_ptr< PeerSet > peerSet_
void addPeers()
Add more peers to the set, if possible.
bool takeTxRootNode(std::string_view data, SHAMapAddNode &san)
Process AS root node received from a peer Call with a lock.
void init(ScopedLockType &collectionLock)
void update(std::uint32_t seq)
bool gotData(std::weak_ptr< Peer >, std::shared_ptr< protocol::TMLedgerData > const &)
Stash a TMLedgerData received from a peer for later processing Returns 'true' if we need to dispatch.
std::vector< NeededHashT > getNeededHashes()
beast::AbstractClock< std::chrono::steady_clock > ClockType
std::vector< UInt256 > neededTxHashes(int max, SHAMapSyncFilter const *filter) const
Rules controlling protocol behavior.
Definition Rules.h:40
bool isInvalid() const
std::string get() const
bool isUseful() const
bool isZero() const
Definition SHAMapHash.h:36
Identifies a node inside a SHAMap.
std::vector< std::pair< SHAMapNodeID, UInt256 > > getMissingNodes(int maxNodes, SHAMapSyncFilter const *filter)
Check for nodes in the SHAMap not available.
SHAMapHash getHash() const
int addRaw(Blob const &vector)
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.
void queueJob(ScopedLockType &)
Queue a job to call invokeOnTimer().
TimeoutCounter(Application &app, UInt256 const &targetHash, std::chrono::milliseconds timeoutInterval, QueueJobParameter &&jobParameter, beast::Journal journal)
bool progress_
Whether forward progress has been made.
beast::Journal journal_
Persistency layer for NodeObject.
Definition Database.h:45
std::shared_ptr< NodeObject > fetchNodeObject(UInt256 const &hash, std::uint32_t ledgerSeq=0, FetchType fetchType=FetchType::Synchronous, bool duplicate=false)
Fetch a node object.
T count_if(T... args)
T emplace_back(T... args)
T for_each(T... args)
T lock(T... args)
T make_shared(T... args)
T make_unique(T... args)
T max(T... args)
@ Array
array value (ordered list)
Definition json_value.h:28
@ Object
object value (collection of name/value pairs).
Definition json_value.h:29
Charge const kFeeMalformedRequest
Schedule of fees charged for imposing load on the server.
Charge const kFeeInvalidData
Use hash_* containers for keys that do not need a cryptographically secure hashing algorithm.
Definition algorithm.h:5
static constexpr auto kReqNodesReply
Number root(Number f, unsigned d)
std::optional< SHAMapNodeID > getSHAMapNodeID(protocol::TMLedgerNode const &ledgerNode, SHAMapTreeNode const &treeNode)
Extracts or reconstructs the SHAMapNodeID from a ledger node proto message.
LedgerHeader deserializeHeader(Slice data, bool hasHash=false)
Deserialize a ledger header from a byte array.
std::string to_string(BaseUInt< Bits, Tag > const &a)
Definition base_uint.h:657
static constexpr std::uint32_t kXrpLedgerEarliestFees
The XRP Ledger mainnet's earliest ledger with a FeeSettings object.
SHAMapTreeNodePtr getTreeNode(std::string_view data)
Deserializes a SHAMapTreeNode from wire format data.
static constexpr auto kMissingNodesFind
Slice makeSlice(std::array< T, N > const &a)
Definition Slice.h:228
static std::vector< UInt256 > neededHashes(UInt256 const &root, SHAMap &map, int max, SHAMapSyncFilter const *filter)
BaseUInt< 256 > UInt256
Definition base_uint.h:580
@ JtLedgerData
Definition Job.h:52
static constexpr auto kLedgerBecomeAggressiveThreshold
static constexpr auto kPeerCountAdd
static constexpr auto kLedgerTimeoutRetriesMax
static constexpr auto kPeerCountStart
@ LedgerMaster
ledger master data for signing
Definition HashPrefix.h:59
constexpr auto kLedgerAcquireTimeout
FeeSettingsEntry< ReadView > FeeSettingsEntryR
std::vector< unsigned char > Blob
Storage for linear binary data.
Definition Blob.h:11
static constexpr auto kReqNodes
LedgerHeader deserializePrefixedHeader(Slice data, bool hasHash=false)
Deserialize a ledger header (prefixed with 4 bytes) from a byte array.
T push_back(T... args)
T reserve(T... args)
T sample(T... args)
T stable_partition(T... args)
T str(T... args)
std::unordered_map< std::shared_ptr< Peer >, int > counts
void update(std::shared_ptr< Peer > &&peer, int dataCount)
void sampleN(std::size_t n, F &&f)
T to_string(T... args)
T unlock(T... args)
T what(T... args)