xrpld
Loading...
Searching...
No Matches
PeerImp.cpp
1#include <xrpld/overlay/detail/PeerImp.h>
2
3#include <xrpld/app/consensus/RCLCxPeerPos.h>
4#include <xrpld/app/consensus/RCLValidations.h>
5#include <xrpld/app/ledger/InboundLedgers.h>
6#include <xrpld/app/ledger/InboundTransactions.h>
7#include <xrpld/app/ledger/LedgerMaster.h>
8#include <xrpld/app/ledger/LedgerNodeHelpers.h>
9#include <xrpld/app/ledger/TransactionMaster.h>
10#include <xrpld/app/misc/Transaction.h>
11#include <xrpld/app/misc/ValidatorList.h>
12#include <xrpld/overlay/Cluster.h>
13#include <xrpld/overlay/ClusterNode.h>
14#include <xrpld/overlay/Peer.h>
15#include <xrpld/overlay/ReduceRelayCommon.h>
16#include <xrpld/overlay/detail/Handshake.h>
17#include <xrpld/overlay/detail/OverlayImpl.h>
18#include <xrpld/overlay/detail/ProtocolMessage.h>
19#include <xrpld/overlay/detail/ProtocolVersion.h>
20#include <xrpld/overlay/detail/TrafficCount.h>
21#include <xrpld/overlay/detail/Tuning.h>
22
23#include <xrpl/basics/Blob.h>
24#include <xrpl/basics/Log.h>
25#include <xrpl/basics/SHAMapHash.h>
26#include <xrpl/basics/Slice.h>
27#include <xrpl/basics/ToString.h>
28#include <xrpl/basics/UptimeClock.h>
29#include <xrpl/basics/base64.h>
30#include <xrpl/basics/base_uint.h>
31#include <xrpl/basics/chrono.h>
32#include <xrpl/basics/random.h>
33#include <xrpl/basics/safe_cast.h>
34#include <xrpl/basics/strHex.h>
35#include <xrpl/beast/utility/Journal.h>
36#include <xrpl/beast/utility/Zero.h>
37#include <xrpl/beast/utility/instrumentation.h>
38#include <xrpl/consensus/Validations.h>
39#include <xrpl/core/HashRouter.h>
40#include <xrpl/core/Job.h>
41#include <xrpl/core/PerfLog.h>
42#include <xrpl/json/json_forwards.h>
43#include <xrpl/json/json_value.h>
44#include <xrpl/ledger/Ledger.h>
45#include <xrpl/peerfinder/Slot.h>
46#include <xrpl/peerfinder/Types.h>
47#include <xrpl/protocol/Feature.h>
48#include <xrpl/protocol/KeyType.h>
49#include <xrpl/protocol/LedgerHeader.h>
50#include <xrpl/protocol/Protocol.h>
51#include <xrpl/protocol/PublicKey.h>
52#include <xrpl/protocol/SField.h>
53#include <xrpl/protocol/STTx.h>
54#include <xrpl/protocol/Serializer.h>
55#include <xrpl/protocol/TxFlags.h>
56#include <xrpl/protocol/digest.h>
57#include <xrpl/protocol/jss.h>
58#include <xrpl/protocol/tokens.h>
59#include <xrpl/resource/Charge.h>
60#include <xrpl/resource/Consumer.h>
61#include <xrpl/resource/Disposition.h>
62#include <xrpl/resource/Fees.h>
63#include <xrpl/resource/Gossip.h>
64#include <xrpl/server/LoadFeeTrack.h>
65#include <xrpl/server/Manifest.h>
66#include <xrpl/server/NetworkOPs.h>
67#include <xrpl/shamap/SHAMap.h>
68#include <xrpl/shamap/SHAMapNodeID.h>
69#include <xrpl/tx/apply.h>
70
71#include <boost/algorithm/string/predicate.hpp>
72#include <boost/asio/bind_executor.hpp>
73#include <boost/asio/buffer.hpp>
74#include <boost/asio/completion_condition.hpp>
75#include <boost/asio/dispatch.hpp>
76#include <boost/asio/error.hpp>
77#include <boost/asio/strand.hpp>
78#include <boost/asio/write.hpp>
79#include <boost/beast/core/multi_buffer.hpp>
80#include <boost/beast/core/ostream.hpp>
81#include <boost/system/system_error.hpp>
82
83#include <google/protobuf/message.h>
84
85#include <xrpl.pb.h>
86
87#include <algorithm>
88#include <atomic>
89#include <chrono>
90#include <cstddef>
91#include <cstdint>
92#include <exception>
93#include <functional>
94#include <map>
95#include <memory>
96#include <mutex>
97#include <numeric>
98#include <optional>
99#include <shared_mutex>
100#include <sstream>
101#include <string>
102#include <string_view>
103#include <tuple>
104#include <utility>
105#include <vector>
106
107using namespace std::chrono_literals;
108
109namespace xrpl {
110
111namespace {
115constexpr std::chrono::milliseconds kPeerHighLatency{300};
116
120constexpr std::chrono::seconds kPeerTimerInterval{60};
121
122} // namespace
123
124// TODO: Remove this exclusion once unit tests are added after the hotfix
125// release.
126
128 Application& app,
129 ID id,
131 HttpRequestType&& request,
132 PublicKey const& publicKey,
134 resource::Consumer consumer,
135 std::unique_ptr<StreamType>&& streamPtr,
136 OverlayImpl& overlay)
137 : Child(overlay)
138 , app_(app)
139 , id_(id)
140 , fingerprint_(getFingerprint(slot->remoteEndpoint(), publicKey, to_string(id)))
142 , sink_(app_.getJournal("Peer"), prefix_)
143 , pSink_(app_.getJournal("Protocol"), prefix_)
144 , journal_(sink_)
146 , streamPtr_(std::move(streamPtr))
147 , socket_(streamPtr_->next_layer().socket())
149 , strand_(boost::asio::make_strand(socket_.get_executor()))
150 , timer_(WaitableTimer{socket_.get_executor()})
151 , remoteAddress_(slot->remoteEndpoint())
152 , overlay_(overlay)
153 , inbound_(true)
154 , protocol_(std::move(protocol))
156 , trackingTime_(ClockType::now())
157 , publicKey_(publicKey)
158 , lastPingTime_(ClockType::now())
159 , creationTime_(ClockType::now())
160 , squelch_(app_.getJournal("Squelch"))
161 , usage_(consumer)
162 , fee_{.fee = resource::kFeeTrivialPeer, .context = ""}
163 , slot_(slot)
164 , request_(std::move(request))
168 ? Compressed::On
169 : Compressed::Off)
171 peerFeatureEnabled(headers_, kFeatureTxrr, app_.config().txReduceRelayEnable))
173 peerFeatureEnabled(headers_, kFeatureLedgerReplay, app_.config().ledgerReplay))
174 , ledgerReplayMsgHandler_(app, app.getLedgerReplayer())
175{
176 JLOG(journal_.info())
177 << "compression enabled " << (compressionEnabled_ == Compressed::On)
178 << " vp reduce-relay base squelch enabled "
179 << peerFeatureEnabled(headers_, kFeatureVprr, app_.config().vpReduceRelayBaseSquelchEnable)
180 << " tx reduce-relay enabled " << txReduceRelayEnabled_;
181}
182
184{
185 bool const inCluster{cluster()};
186
187 overlay_.deletePeer(id_);
188 overlay_.onPeerDeactivate(id_);
189 overlay_.peerFinder().onClosed(slot_);
190 overlay_.remove(slot_);
191
192 if (inCluster)
193 {
194 JLOG(journal_.warn()) << name() << " left cluster";
195 }
196}
197
198// Helper function to check for valid UInt256 values in protobuf buffers
199static bool
201{
202 return pBuffStr.size() == UInt256::size();
203}
204
205void
207{
208 dispatch(strand_, [self = shared_from_this()]() {
209 auto parseLedgerHash = [](std::string_view value) -> std::optional<UInt256> {
210 if (UInt256 ret; ret.parseHex(value))
211 return ret;
212
213 if (auto const s = base64Decode(value); s.size() == UInt256::size())
214 return UInt256::fromRaw(s);
215
216 return std::nullopt;
217 };
218
220 std::optional<UInt256> previous;
221
222 if (auto const iter = self->headers_.find("Closed-Ledger"); iter != self->headers_.end())
223 {
224 closed = parseLedgerHash(iter->value());
225
226 if (!closed)
227 self->fail("Malformed handshake data (1)");
228 }
229
230 if (auto const iter = self->headers_.find("Previous-Ledger"); iter != self->headers_.end())
231 {
232 previous = parseLedgerHash(iter->value());
233
234 if (!previous)
235 self->fail("Malformed handshake data (2)");
236 }
237
238 if (previous && !closed)
239 self->fail("Malformed handshake data (3)");
240
241 {
242 std::scoped_lock const sl(self->recentLock_);
243 if (closed)
244 self->closedLedgerHash_ = *closed;
245 if (previous)
246 self->previousLedgerHash_ = *previous;
247 }
248
249 if (self->inbound_)
250 {
251 self->doAccept();
252 }
253 else
254 {
255 self->doProtocolStart();
256 }
257
258 // Anything else that needs to be done with the connection should be
259 // done in doProtocolStart
260 });
261}
262
263void
265{
266 dispatch(strand_, [self = shared_from_this()]() {
267 if (!self->socket_.is_open())
268 return;
269
270 self->close();
271 });
272}
273
274//------------------------------------------------------------------------------
275
276void
278{
279 dispatch(strand_, [self = shared_from_this(), m]() {
280 if (self->gracefulClose_)
281 return;
282 if (self->detaching_)
283 return;
284 if (!self->socket_.is_open())
285 return;
286
287 auto validator = m->getValidatorKey();
288 if (validator && !self->squelch_.expireSquelch(*validator))
289 {
290 self->overlay_.reportOutboundTraffic(
292 static_cast<int>(m->getBuffer(self->compressionEnabled_).size()));
293 return;
294 }
295
296 // report categorized outgoing traffic
297 self->overlay_.reportOutboundTraffic(
298 safeCast<TrafficCount::Category>(m->getCategory()),
299 static_cast<int>(m->getBuffer(self->compressionEnabled_).size()));
300
301 // report total outgoing traffic
302 self->overlay_.reportOutboundTraffic(
304 static_cast<int>(m->getBuffer(self->compressionEnabled_).size()));
305
306 auto sendqSize = self->sendQueue_.size();
307
308 if (sendqSize < tuning::kTargetSendQueue)
309 {
310 // To detect a peer that does not read from their
311 // side of the connection, we expect a peer to have
312 // a small sendq periodically
313 self->largeSendq_ = 0;
314 }
315 else if (
316 auto sink = self->journal_.debug();
317 sink && (sendqSize % tuning::kSendQueueLogFreq) == 0)
318 {
319 std::string const n = self->name();
320 sink << n << " sendq: " << sendqSize;
321 }
322
323 self->sendQueue_.push(m);
324
325 if (sendqSize != 0)
326 return;
327
328 boost::asio::async_write(
329 self->stream_,
330 boost::asio::buffer(self->sendQueue_.front()->getBuffer(self->compressionEnabled_)),
331 bind_executor(self->strand_, [self](ErrorCode const& ec, std::size_t bytesTransferred) {
332 self->onWriteMessage(ec, bytesTransferred);
333 }));
334 });
335}
336
337void
339{
340 dispatch(strand_, [self = shared_from_this()]() {
341 if (!self->txQueue_.empty())
342 {
343 protocol::TMHaveTransactions ht;
345 self->txQueue_, [&](auto const& hash) { ht.add_hashes(hash.data(), hash.size()); });
346 JLOG(self->pJournal_.trace()) << "sendTxQueue " << self->txQueue_.size();
347 self->txQueue_.clear();
348 self->send(std::make_shared<Message>(ht, protocol::mtHAVE_TRANSACTIONS));
349 }
350 });
351}
352
353void
355{
356 dispatch(strand_, [self = shared_from_this(), hash]() {
357 if (self->txQueue_.size() == reduce_relay::kMaxTxQueueSize)
358 {
359 JLOG(self->pJournal_.warn()) << "addTxQueue exceeds the cap";
360 self->sendTxQueue();
361 }
362
363 self->txQueue_.insert(hash);
364 JLOG(self->pJournal_.trace()) << "addTxQueue " << self->txQueue_.size();
365 });
366}
367
368void
370{
371 dispatch(strand_, [self = shared_from_this(), hash]() {
372 auto removed = self->txQueue_.erase(hash);
373 JLOG(self->pJournal_.trace()) << "removeTxQueue " << removed;
374 });
375}
376
377void
379{
380 dispatch(strand_, [self = shared_from_this(), fee, context]() {
381 if ((self->usage_.charge(fee, context) == resource::Disposition::Drop) &&
382 self->usage_.disconnect(self->pJournal_))
383 {
384 // Idempotent: only the first worker to observe Drop counts the
385 // metric and posts fail(). Without the guard, several queued
386 // workers can all see Drop before fail() lands on the strand,
387 // overcounting peerDisconnectsCharges_ and posting duplicate
388 // shutdowns. fail(std::string const&) self-posts to strand_
389 // when invoked off-strand.
390 bool expected = false;
391 if (self->chargeDisconnectFired_.compare_exchange_strong(
392 expected, true, std::memory_order_acq_rel))
393 {
394 self->overlay_.incPeerDisconnectCharges();
395 self->fail("charge: Resources");
396 }
397 }
398 });
399}
400
401//------------------------------------------------------------------------------
402
403bool
405{
406 auto const iter = headers_.find("Crawl");
407 if (iter == headers_.end())
408 return false;
409 return boost::iequals(iter->value(), "public");
410}
411
412bool
414{
415 return app_.getCluster().isMember(publicKey_);
416}
417
420{
421 if (inbound_)
422 return headers_["User-Agent"];
423 return headers_["Server"];
424}
425
428{
430
431 ret[jss::public_key] = toBase58(TokenType::NodePublic, publicKey_);
432 ret[jss::address] = remoteAddress_.toString();
433
434 if (inbound_)
435 ret[jss::inbound] = true;
436
437 if (cluster())
438 {
439 ret[jss::cluster] = true;
440
441 if (auto const n = name(); !n.empty())
442 {
443 // Could move here if json::Value supported moving from a string
444 ret[jss::name] = n;
445 }
446 }
447
448 if (auto const d = domain(); !d.empty())
449 ret[jss::server_domain] = std::string{d};
450
451 if (auto const nid = headers_["Network-ID"]; !nid.empty())
452 ret[jss::network_id] = std::string{nid};
453
454 ret[jss::load] = usage_.balance();
455
456 if (auto const version = getVersion(); !version.empty())
457 ret[jss::version] = std::string{version};
458
459 ret[jss::protocol] = to_string(protocol_);
460
461 {
463 if (latency_)
464 ret[jss::latency] = static_cast<json::UInt>(latency_->count());
465 }
466
467 ret[jss::uptime] =
469
470 std::uint32_t minSeq = 0, maxSeq = 0;
471 ledgerRange(minSeq, maxSeq);
472
473 if ((minSeq != 0) || (maxSeq != 0))
474 ret[jss::complete_ledgers] = std::to_string(minSeq) + " - " + std::to_string(maxSeq);
475
476 switch (tracking_.load())
477 {
479 ret[jss::track] = "diverged";
480 break;
481
483 ret[jss::track] = "unknown";
484 break;
485
487 // Nothing to do here
488 break;
489 }
490
491 UInt256 closedLedgerHash;
492 protocol::TMStatusChange lastStatus;
493 {
495 closedLedgerHash = closedLedgerHash_;
496 lastStatus = lastStatus_;
497 }
498
499 if (closedLedgerHash != beast::kZero)
500 ret[jss::ledger] = to_string(closedLedgerHash);
501
502 if (lastStatus.has_newstatus())
503 {
504 switch (lastStatus.newstatus())
505 {
506 case protocol::nsCONNECTING:
507 ret[jss::status] = "connecting";
508 break;
509
510 case protocol::nsCONNECTED:
511 ret[jss::status] = "connected";
512 break;
513
514 case protocol::nsMONITORING:
515 ret[jss::status] = "monitoring";
516 break;
517
518 case protocol::nsVALIDATING:
519 ret[jss::status] = "validating";
520 break;
521
522 case protocol::nsSHUTTING:
523 ret[jss::status] = "shutting";
524 break;
525
526 default:
527 JLOG(pJournal_.warn()) << "Unknown status: " << lastStatus.newstatus();
528 }
529 }
530
531 ret[jss::metrics] = json::Value(json::ValueType::Object);
532 ret[jss::metrics][jss::total_bytes_recv] = std::to_string(metrics_.recv.totalBytes());
533 ret[jss::metrics][jss::total_bytes_sent] = std::to_string(metrics_.sent.totalBytes());
534 ret[jss::metrics][jss::avg_bps_recv] = std::to_string(metrics_.recv.averageBytes());
535 ret[jss::metrics][jss::avg_bps_sent] = std::to_string(metrics_.sent.averageBytes());
536
537 return ret;
538}
539
540bool
542{
543 switch (f)
544 {
546 return protocol_ >= makeProtocol(2, 3);
549 }
550 return false;
551}
552
553//------------------------------------------------------------------------------
554
555bool
557{
558 {
560 if ((seq != 0) && (seq >= minLedger_) && (seq <= maxLedger_) &&
561 (tracking_.load() == Tracking::Converged))
562 return true;
564 return true;
565 }
566 return false;
567}
568
569void
571{
573
574 minSeq = minLedger_;
575 maxSeq = maxLedger_;
576}
577
578bool
579PeerImp::hasTxSet(UInt256 const& hash) const
580{
582 return std::ranges::find(recentTxSets_, hash) != recentTxSets_.end();
583}
584
585void
587{
588 // Operations on closedLedgerHash_ and previousLedgerHash_ must be
589 // guarded by recentLock_.
592 closedLedgerHash_.zero();
593}
594
595bool
597{
599 return (tracking_ != Tracking::Diverged) && (uMin >= minLedger_) && (uMax <= maxLedger_);
600}
601
602//------------------------------------------------------------------------------
603
604void
606{
607 XRPL_ASSERT(strand_.running_in_this_thread(), "xrpl::PeerImp::close : strand in this thread");
608 if (!socket_.is_open())
609 return;
610
611 detaching_ = true; // DEPRECATED
612
613 cancelTimer();
614 ErrorCode ec;
615 socket_.close(ec); // NOLINT(bugprone-unused-return-value)
616
617 overlay_.incPeerDisconnect();
618 JLOG((inbound_ ? journal_.debug() : journal_.info())) << "close: Closed";
619}
620
621void
623{
624 dispatch(strand_, [self = shared_from_this(), reason]() {
625 if (self->journal_.active(beast::Severity::Warning) && self->socket_.is_open())
626 {
627 std::string const n = self->name();
628 JLOG(self->journal_.warn()) << n << " failed: " << reason;
629 }
630 self->close();
631 });
632}
633
634void
636{
637 XRPL_ASSERT(strand_.running_in_this_thread(), "xrpl::PeerImp::fail : strand in this thread");
638 if (!socket_.is_open())
639 return;
640
641 JLOG(journal_.warn()) << name << ": " << ec.message();
642
643 close();
644}
645
646void
648{
649 XRPL_ASSERT(
650 strand_.running_in_this_thread(), "xrpl::PeerImp::gracefulClose : strand in this thread");
651 XRPL_ASSERT(socket_.is_open(), "xrpl::PeerImp::gracefulClose : socket is open");
652 XRPL_ASSERT(!gracefulClose_, "xrpl::PeerImp::gracefulClose : socket is not closing");
653 gracefulClose_ = true;
654 if (!sendQueue_.empty())
655 return;
656 setTimer();
657 stream_.async_shutdown(bind_executor(
658 strand_, [self = shared_from_this()](ErrorCode const& ec) { self->onShutdown(ec); }));
659}
660
661void
663{
664 try
665 {
666 timer_.expires_after(kPeerTimerInterval);
667 }
668 catch (boost::system::system_error const& e)
669 {
670 JLOG(journal_.error()) << "setTimer: " << e.code();
671 return;
672 }
673 timer_.async_wait(bind_executor(
674 strand_, [self = shared_from_this()](ErrorCode const& ec) { self->onTimer(ec); }));
675}
676
677// convenience for ignoring the error code
678void
680{
681 try
682 {
683 timer_.cancel();
684 }
685 catch (boost::system::system_error const&) // NOLINT(bugprone-empty-catch)
686 {
687 // ignored
688 }
689}
690
691//------------------------------------------------------------------------------
692
695{
697 ss << "[" << fingerprint << "] ";
698 return ss.str();
699}
700
701void
703{
704 if (!socket_.is_open())
705 return;
706
707 if (ec)
708 {
709 if (ec == boost::asio::error::operation_aborted)
710 return;
711
712 // This should never happen
713 JLOG(journal_.error()) << "onTimer: " << ec.message();
714 close();
715 return;
716 }
717
719 {
720 fail("Large send queue");
721 return;
722 }
723
724 if (auto const t = tracking_.load(); !inbound_ && t != Tracking::Converged)
725 {
726 ClockType::duration duration;
727
728 {
730 duration = ClockType::now() - trackingTime_;
731 }
732
733 if ((t == Tracking::Diverged && (duration > app_.config().maxDivergedTime)) ||
734 (t == Tracking::Unknown && (duration > app_.config().maxUnknownTime)))
735 {
736 overlay_.peerFinder().onFailure(slot_);
737 fail("Not useful");
738 return;
739 }
740 }
741
742 // Already waiting for PONG
743 if (lastPingSeq_)
744 {
745 fail("Ping Timeout");
746 return;
747 }
748
751
752 protocol::TMPing message;
753 message.set_type(protocol::TMPing::ptPING);
754 message.set_seq(*lastPingSeq_);
755
756 send(std::make_shared<Message>(message, protocol::mtPING));
757
758 setTimer();
759}
760
761void
763{
764 cancelTimer();
765
766 if (ec)
767 {
768 // - eof: the stream was cleanly closed
769 // - operation_aborted: an expired timer (slow shutdown)
770 // - stream_truncated: the tcp connection closed (no handshake) it could
771 // occur if a peer does not perform a graceful disconnect
772 // - broken_pipe: the peer is gone
773 bool const shouldLog =
774 (ec != boost::asio::error::eof && ec != boost::asio::error::operation_aborted &&
775 !ec.message().contains("application data after close notify"));
776
777 if (shouldLog)
778 {
779 JLOG(journal_.debug()) << "onShutdown: " << ec.message();
780 }
781 }
782
783 close();
784}
785
786//------------------------------------------------------------------------------
787void
789{
790 XRPL_ASSERT(readBuffer_.size() == 0, "xrpl::PeerImp::doAccept : empty read buffer");
791
792 auto const sharedValue = makeSharedValue(*streamPtr_, journal_);
793
794 // This shouldn't fail since we already computed
795 // the shared value successfully in OverlayImpl
796 if (!sharedValue)
797 {
798 fail("makeSharedValue: Unexpected failure");
799 return;
800 }
801
802 JLOG(journal_.info()) << "Protocol: " << to_string(protocol_);
803
804 if (auto member = app_.getCluster().member(publicKey_))
805 {
806 {
807 std::unique_lock const lock{nameMutex_};
808 name_ = *member;
809 }
810 JLOG(journal_.info()) << "Cluster name: " << *member;
811 }
812
813 overlay_.activate(shared_from_this());
814
815 // XXX Set timer: connection is in grace period to be useful.
816 // XXX Set timer: connection idle (idle may vary depending on connection
817 // type.)
818
820
821 boost::beast::ostream(*writeBuffer) << makeResponse(
822 !overlay_.peerFinder().config().peerPrivate,
823 request_,
824 overlay_.setup().publicIp,
825 remoteAddress_.address(),
826 *sharedValue,
827 overlay_.setup().networkID,
828 protocol_,
829 app_);
830
831 // Write the whole buffer and only start protocol when that's done.
832 boost::asio::async_write(
833 stream_,
834 writeBuffer->data(),
835 boost::asio::transfer_all(),
836 bind_executor(
837 strand_,
838 [this, writeBuffer, self = shared_from_this()](
839 ErrorCode ec, std::size_t bytesTransferred) {
840 if (!socket_.is_open())
841 return;
842 if (ec)
843 {
844 if (ec == boost::asio::error::operation_aborted)
845 return;
846
847 fail("onWriteResponse", ec);
848 return;
849 }
850
851 if (writeBuffer->size() == bytesTransferred)
852 {
854 return;
855 }
856 fail("Failed to write header");
857 return;
858 }));
859}
860
863{
864 std::shared_lock const readLock{nameMutex_};
865 return name_;
866}
867
870{
871 return headers_["Server-Domain"];
872}
873
874//------------------------------------------------------------------------------
875
876// Protocol logic
877
878void
880{
882
883 // Send all the validator lists that have been loaded
884 if (inbound_)
885 {
886 app_.getValidators().forEachAvailable(
887 [&](std::string const& manifest,
888 std::uint32_t version,
890 PublicKey const& pubKey,
891 std::size_t maxSequence,
892 UInt256 const& hash) {
894 *this,
895 0,
896 pubKey,
897 maxSequence,
898 version,
899 manifest,
900 blobInfos,
901 app_.getHashRouter(),
902 pJournal_);
903
904 // Don't send it next time.
905 app_.getHashRouter().addSuppressionPeer(hash, id_);
906 });
907 }
908
909 if (auto m = overlay_.getManifestsMessage())
910 send(m);
911
912 setTimer();
913}
914
915// Called repeatedly with protocol message data
916void
918{
919 if (!socket_.is_open())
920 return;
921
922 if (ec)
923 {
924 if (ec == boost::asio::error::operation_aborted)
925 return;
926
927 if (ec == boost::asio::error::eof)
928 {
929 JLOG(journal_.info()) << "EOF";
931 return;
932 }
933
934 fail("onReadMessage", ec);
935 return;
936 }
937
938 if (auto stream = journal_.trace())
939 {
940 stream << "onReadMessage: "
941 << (bytesTransferred > 0 ? to_string(bytesTransferred) + " bytes" : "");
942 }
943
944 metrics_.recv.addMessage(bytesTransferred);
945
946 readBuffer_.commit(bytesTransferred);
947
948 auto hint = tuning::kReadBufferBytes;
949
950 while (readBuffer_.size() > 0)
951 {
952 std::size_t bytesConsumed = 0;
953
954 using namespace std::chrono_literals;
955 std::tie(bytesConsumed, ec) = perf::measureDurationAndLog(
956 [&]() { return invokeProtocolMessage(readBuffer_.data(), *this, hint); },
957 "invokeProtocolMessage",
958 350ms,
959 journal_);
960
961 if (ec)
962 {
963 fail("onReadMessage", ec);
964 return;
965 }
966
967 if (!socket_.is_open())
968 return;
969
970 if (gracefulClose_)
971 return;
972
973 if (bytesConsumed == 0)
974 break;
975 readBuffer_.consume(bytesConsumed);
976 }
977
978 // Timeout on writes only
979 stream_.async_read_some(
981 bind_executor(
982 strand_,
983 [self = shared_from_this()](ErrorCode const& ec, std::size_t bytesTransferred) {
984 self->onReadMessage(ec, bytesTransferred);
985 }));
986}
987
988void
990{
991 if (!socket_.is_open())
992 return;
993
994 if (ec)
995 {
996 if (ec == boost::asio::error::operation_aborted)
997 return;
998
999 fail("onWriteMessage", ec);
1000 return;
1001 }
1002 if (auto stream = journal_.trace())
1003 {
1004 stream << "onWriteMessage: "
1005 << (bytesTransferred > 0 ? to_string(bytesTransferred) + " bytes" : "");
1006 }
1007
1008 metrics_.sent.addMessage(bytesTransferred);
1009
1010 XRPL_ASSERT(!sendQueue_.empty(), "xrpl::PeerImp::onWriteMessage : non-empty send buffer");
1011 sendQueue_.pop();
1012 if (!sendQueue_.empty())
1013 {
1014 // Timeout on writes only
1015 boost::asio::async_write(
1016 stream_,
1017 boost::asio::buffer(sendQueue_.front()->getBuffer(compressionEnabled_)),
1018 bind_executor(
1019 strand_,
1020 [self = shared_from_this()](ErrorCode const& ec, std::size_t bytesTransferred) {
1021 self->onWriteMessage(ec, bytesTransferred);
1022 }));
1023 return;
1024 }
1025
1026 if (gracefulClose_)
1027 {
1028 stream_.async_shutdown(bind_executor(
1029 strand_, [self = shared_from_this()](ErrorCode const& ec) { self->onShutdown(ec); }));
1030 return;
1031 }
1032}
1033
1034//------------------------------------------------------------------------------
1035//
1036// ProtocolHandler
1037//
1038//------------------------------------------------------------------------------
1039
1040void
1042{
1043 // TODO
1044}
1045
1046void
1048 std::uint16_t type,
1050 std::size_t size,
1051 std::size_t uncompressedSize,
1052 bool isCompressed)
1053{
1054 auto const name = protocolMessageName(type);
1055 loadEvent_ = app_.getJobQueue().makeLoadEvent(JtPeer, name);
1056 fee_ = {.fee = resource::kFeeTrivialPeer, .context = name};
1057
1058 auto const category = TrafficCount::attribute(
1059 TrafficCount::categorize(*m, static_cast<protocol::MessageType>(type), true),
1061
1062 // report total incoming traffic
1063 overlay_.reportInboundTraffic(TrafficCount::Category::Total, static_cast<int>(size));
1064
1065 // increase the traffic received for a specific category
1066 overlay_.reportInboundTraffic(category, static_cast<int>(size));
1067
1068 using namespace protocol;
1069 if ((type == MessageType::mtTRANSACTION || type == MessageType::mtHAVE_TRANSACTIONS ||
1070 type == MessageType::mtTRANSACTIONS ||
1071 // GET_OBJECTS
1073 // GET_LEDGER
1076 // LEDGER_DATA
1078 category == TrafficCount::Category::GlTscGet) &&
1079 (txReduceRelayEnabled() || app_.config().txReduceRelayMetrics))
1080 {
1081 overlay_.addTxMetrics(static_cast<MessageType>(type), static_cast<std::uint64_t>(size));
1082 }
1083 JLOG(journal_.trace()) << "onMessageBegin: " << type << " " << size << " " << uncompressedSize
1084 << " " << isCompressed;
1085}
1086
1087void
1093
1094void
1096{
1097 auto const s = m->list_size();
1098
1099 if (s == 0)
1100 {
1101 fee_.update(resource::kFeeUselessData, "empty");
1102 return;
1103 }
1104
1105 if (s > 100)
1106 fee_.update(resource::kFeeModerateBurdenPeer, "oversize");
1107
1108 // OverlayImpl::onManifests bounds the untrusted work and charges the fee
1109 // if the untrusted count exceeds the per-message cap; trusted manifests
1110 // are always processed and not counted against it.
1111 app_.getJobQueue().addJob(JtManifest, "RcvManifests", [this, that = shared_from_this(), m]() {
1112 overlay_.onManifests(m, that);
1113 });
1114}
1115
1116void
1118{
1119 if (m->type() == protocol::TMPing::ptPING)
1120 {
1121 // We have received a ping request, reply with a pong.
1122 fee_.update(resource::kFeeModerateBurdenPeer, "ping request");
1123 protocol::TMPing pong;
1124 pong.set_type(protocol::TMPing::ptPONG);
1125 if (m->has_seq())
1126 pong.set_seq(m->seq());
1127 send(std::make_shared<Message>(pong, protocol::mtPING));
1128 return;
1129 }
1130
1131 if (m->type() == protocol::TMPing::ptPONG && m->has_seq())
1132 {
1133 // Only reset the ping sequence if we actually received a
1134 // PONG with the correct cookie. That way, any peers which
1135 // respond with incorrect cookies will eventually time out.
1136 if (m->seq() == lastPingSeq_)
1137 {
1138 lastPingSeq_.reset();
1139
1140 // Update latency estimate
1141 auto const rtt =
1142 std::chrono::round<std::chrono::milliseconds>(ClockType::now() - lastPingTime_);
1143
1144 std::scoped_lock const sl(recentLock_);
1145
1146 if (latency_)
1147 {
1148 latency_ = (*latency_ * 7 + rtt) / 8;
1149 }
1150 else
1151 {
1152 latency_ = rtt;
1153 }
1154 }
1155
1156 return;
1157 }
1158}
1159
1160void
1162{
1163 // VFALCO NOTE I think we should drop the peer immediately
1164 if (!cluster())
1165 {
1166 fee_.update(resource::kFeeUselessData, "unknown cluster");
1167 return;
1168 }
1169
1170 for (int i = 0; i < m->clusternodes().size(); ++i)
1171 {
1172 protocol::TMClusterNode const& node = m->clusternodes(i);
1173
1175 if (node.has_nodename())
1176 name = node.nodename();
1177
1178 auto const publicKey = parseBase58<PublicKey>(TokenType::NodePublic, node.publickey());
1179
1180 // NIKB NOTE We should drop the peer immediately if
1181 // they send us a public key we can't parse
1182 if (publicKey)
1183 {
1184 auto const reportTime = NetClock::time_point{NetClock::duration{node.reporttime()}};
1185
1186 app_.getCluster().update(*publicKey, name, node.nodeload(), reportTime);
1187 }
1188 }
1189
1190 int const loadSources = m->loadsources().size();
1191 if (loadSources != 0)
1192 {
1193 resource::Gossip gossip;
1194 gossip.items.reserve(loadSources);
1195 for (int i = 0; i < m->loadsources().size(); ++i)
1196 {
1197 protocol::TMLoadSource const& node = m->loadsources(i);
1199 item.address = beast::ip::Endpoint::fromString(node.name());
1200 item.balance = node.cost();
1201 if (item.address != beast::ip::Endpoint())
1202 gossip.items.push_back(item);
1203 }
1204 overlay_.resourceManager().importConsumers(name(), gossip);
1205 }
1206
1207 // Calculate the cluster fee:
1208 auto const thresh = app_.getTimeKeeper().now() - 90s;
1209 std::uint32_t clusterFee = 0;
1210
1212 fees.reserve(app_.getCluster().size());
1213
1214 app_.getCluster().forEach([&fees, thresh](ClusterNode const& status) {
1215 if (status.getReportTime() >= thresh)
1216 fees.push_back(status.getLoadFee());
1217 });
1218
1219 if (!fees.empty())
1220 {
1221 auto const index = fees.size() / 2;
1222 std::nth_element(fees.begin(), fees.begin() + index, fees.end());
1223 clusterFee = fees[index];
1224 }
1225
1226 app_.getFeeTrack().setClusterFee(clusterFee);
1227}
1228
1229void
1231{
1232 // Don't allow endpoints from peers that are not known tracking or are
1233 // not using a version of the message that we support:
1234 if (tracking_.load() != Tracking::Converged || m->version() != 2)
1235 return;
1236
1237 // The number is arbitrary and doesn't have any real significance or
1238 // implication for the protocol.
1239 if (m->endpoints_v2().size() >= 1024)
1240 {
1241 fee_.update(resource::kFeeUselessData, "endpoints too large");
1242 return;
1243 }
1244
1246 endpoints.reserve(m->endpoints_v2().size());
1247
1248 auto malformed = 0;
1249 for (auto const& tm : m->endpoints_v2())
1250 {
1251 auto result = beast::ip::Endpoint::fromStringChecked(tm.endpoint());
1252
1253 if (!result)
1254 {
1255 JLOG(pJournal_.error())
1256 << "failed to parse incoming endpoint: {" << tm.endpoint() << "}";
1257 malformed++;
1258 continue;
1259 }
1260
1261 // If hops == 0, this Endpoint describes the peer we are connected
1262 // to -- in that case, we take the remote address seen on the
1263 // socket and store that in the ip::Endpoint. If this is the first
1264 // time, then we'll verify that their listener can receive incoming
1265 // by performing a connectivity test. if hops > 0, then we just
1266 // take the address/port we were given
1267 if (tm.hops() == 0)
1268 result = remoteAddress_.atPort(result->port());
1269
1270 endpoints.emplace_back(*result, tm.hops());
1271 }
1272
1273 // Charge the peer for each malformed endpoint. As there still may be
1274 // multiple valid endpoints we don't return early.
1275 if (malformed > 0)
1276 {
1277 fee_.update(
1278 resource::kFeeInvalidData * malformed,
1279 std::to_string(malformed) + " malformed endpoints");
1280 }
1281
1282 if (!endpoints.empty())
1283 overlay_.peerFinder().onEndpoints(slot_, endpoints);
1284}
1285
1286void
1291
1292void
1295 bool eraseTxQueue,
1296 bool batch)
1297{
1298 XRPL_ASSERT(eraseTxQueue != batch, ("xrpl::PeerImp::handleTransaction : valid inputs"));
1299 if (tracking_.load() == Tracking::Diverged)
1300 return;
1301
1302 if (app_.getOPs().isNeedNetworkLedger())
1303 {
1304 // If we've never been in synch, there's nothing we can do
1305 // with a transaction
1306 JLOG(pJournal_.debug()) << "Ignoring incoming transaction: Need network ledger";
1307 return;
1308 }
1309
1310 SerialIter sit(makeSlice(m->rawtransaction()));
1311
1312 try
1313 {
1314 auto stx = std::make_shared<STTx const>(sit);
1315 UInt256 const txID = stx->getTransactionID();
1316
1317 // Charge strongly for attempting to relay a txn with tfInnerBatchTxn
1318 // LCOV_EXCL_START
1319 /*
1320 There is no need to check whether the featureBatchV1_1 amendment is
1321 enabled.
1322
1323 * If the `tfInnerBatchTxn` flag is set, and the amendment is
1324 enabled, then it's an invalid transaction because inner batch
1325 transactions should not be relayed.
1326 * If the `tfInnerBatchTxn` flag is set, and the amendment is *not*
1327 enabled, then the transaction is malformed because it's using an
1328 "unknown" flag. There's no need to waste the resources to send it
1329 to the transaction engine.
1330
1331 We don't normally check transaction validity at this level, but
1332 since we _need_ to check it when the amendment is enabled, we may as
1333 well drop it if the flag is set regardless.
1334 */
1335 if (stx->isFlag(tfInnerBatchTxn))
1336 {
1337 JLOG(pJournal_.warn()) << "Ignoring Network relayed Tx containing "
1338 "tfInnerBatchTxn (handleTransaction).";
1339 fee_.update(resource::kFeeModerateBurdenPeer, "inner batch txn");
1340 return;
1341 }
1342 // LCOV_EXCL_STOP
1343
1345 static constexpr std::chrono::seconds kTxInterval = 10s;
1346
1347 if (!app_.getHashRouter().shouldProcess(txID, id_, flags, kTxInterval))
1348 {
1349 // we have seen this transaction recently
1350 if (any(flags & HashRouterFlags::BAD))
1351 {
1352 fee_.update(resource::kFeeUselessData, "known bad");
1353 JLOG(pJournal_.debug()) << "Ignoring known bad tx " << txID;
1354 }
1355
1356 // Erase only if the server has seen this tx. If the server has not
1357 // seen this tx then the tx could not has been queued for this peer.
1358 else if (eraseTxQueue && txReduceRelayEnabled())
1359 {
1360 removeTxQueue(txID);
1361 }
1362
1363 overlay_.reportInboundTraffic(
1365
1366 return;
1367 }
1368
1369 JLOG(pJournal_.debug()) << "Got tx " << txID;
1370
1371 bool checkSignature = true;
1372 if (cluster())
1373 {
1374 if (!m->has_deferred() || !m->deferred())
1375 {
1376 // Skip local checks if a server we trust
1377 // put the transaction in its open ledger
1378 flags |= HashRouterFlags::TRUSTED;
1379 }
1380
1381 // for non-validator nodes only -- localPublicKey is set for
1382 // validators only
1383 if (!app_.getValidationPublicKey())
1384 {
1385 // For now, be paranoid and have each validator
1386 // check each transaction, regardless of source
1387 checkSignature = false;
1388 }
1389 }
1390
1391 if (app_.getLedgerMaster().getValidatedLedgerAge() > 4min)
1392 {
1393 JLOG(pJournal_.trace()) << "No new transactions until synchronized";
1394 }
1395 else if (app_.getJobQueue().getJobCount(JtTransaction) > app_.config().maxTransactions)
1396 {
1397 overlay_.incJqTransOverflow();
1398 JLOG(pJournal_.info()) << "Transaction queue is full";
1399 }
1400 else
1401 {
1402 app_.getJobQueue().addJob(
1404 "RcvCheckTx",
1406 flags,
1407 checkSignature,
1408 batch,
1409 stx]() {
1410 if (auto peer = weak.lock())
1411 peer->checkTransaction(flags, checkSignature, stx, batch);
1412 });
1413 }
1414 }
1415 catch (std::exception const& ex)
1416 {
1418 {
1419 fee_.update(resource::kFeeInvalidData, "tx invalid");
1420 }
1421 JLOG(pJournal_.warn()) << "Transaction invalid: " << strHex(m->rawtransaction())
1422 << ". Exception: " << ex.what();
1423 }
1424}
1425
1426void
1428{
1429 auto badData = [&](std::string const& msg) {
1430 fee_.update(resource::kFeeInvalidData, "get_ledger " + msg);
1431 JLOG(pJournal_.warn()) << "TMGetLedger: " << msg;
1432 };
1433 auto const itype{m->itype()};
1434
1435 // Verify ledger info type
1436 if (itype < protocol::liBASE || itype > protocol::liTS_CANDIDATE)
1437 {
1438 badData("Invalid ledger info type");
1439 return;
1440 }
1441
1442 auto const ltype = [&m]() -> std::optional<::protocol::TMLedgerType> {
1443 if (m->has_ltype())
1444 return m->ltype();
1445 return std::nullopt;
1446 }();
1447
1448 if (itype == protocol::liTS_CANDIDATE)
1449 {
1450 if (!m->has_ledgerhash())
1451 {
1452 badData("Invalid TX candidate set, missing TX set hash");
1453 return;
1454 }
1455 }
1456 else if (
1457 !m->has_ledgerhash() && !m->has_ledgerseq() && (!ltype || *ltype != protocol::ltCLOSED))
1458 {
1459 badData("Invalid request");
1460 return;
1461 }
1462
1463 // Verify ledger type
1464 if (ltype && (*ltype < protocol::ltACCEPTED || *ltype > protocol::ltCLOSED))
1465 {
1466 badData("Invalid ledger type");
1467 return;
1468 }
1469
1470 // Verify ledger hash
1471 if (m->has_ledgerhash() && !stringIsUInt256Sized(m->ledgerhash()))
1472 {
1473 badData("Invalid ledger hash");
1474 return;
1475 }
1476
1477 // Verify ledger sequence
1478 if (m->has_ledgerseq())
1479 {
1480 auto const ledgerSeq{m->ledgerseq()};
1481
1482 // Check if within a reasonable range
1483 using namespace std::chrono_literals;
1484 if (app_.getLedgerMaster().getValidatedLedgerAge() <= 10s &&
1485 ledgerSeq > app_.getLedgerMaster().getValidLedgerIndex() + 10)
1486 {
1487 badData("Invalid ledger sequence " + std::to_string(ledgerSeq));
1488 return;
1489 }
1490 }
1491
1492 // Verify ledger node counts. Full parsing of the node IDs is deferred to the job, so the I/O
1493 // thread is not burdened with SHAMapNodeID deserialization for every TMGetLedger message.
1494 if (itype != protocol::liBASE)
1495 {
1496 if (m->nodeids_size() <= 0)
1497 {
1498 badData("Invalid ledger node IDs");
1499 return;
1500 }
1501
1502 if (m->nodeids_size() > tuning::kHardMaxReplyNodes)
1503 {
1504 badData(
1505 "Requested number of ledger node IDs must be less than or equal to " +
1507 return;
1508 }
1509 }
1510
1511 // Verify query type
1512 if (m->has_querytype() && m->querytype() != protocol::qtINDIRECT)
1513 {
1514 badData("Invalid query type");
1515 return;
1516 }
1517
1518 // Verify query depth
1519 if (m->has_querydepth())
1520 {
1521 if (m->querydepth() > tuning::kMaxQueryDepth || itype == protocol::liBASE)
1522 {
1523 badData("Invalid query depth");
1524 return;
1525 }
1526 }
1527
1528 // Queue a job to process the request.
1530 app_.getJobQueue().addJob(JtLedgerReq, "RcvGetLedger", [weak, m, itype]() {
1531 auto peer = weak.lock();
1532 if (!peer)
1533 return;
1534
1536 bool tooManyNodeIds = false;
1537 if (itype != protocol::liBASE)
1538 {
1539 nodeIDs.reserve(std::min(m->nodeids_size(), tuning::kSoftMaxReplyNodes));
1540 for (auto const& nodeId : m->nodeids())
1541 {
1542 if (nodeIDs.size() >= tuning::kSoftMaxReplyNodes)
1543 {
1544 // The peer requested too many node IDs. Continue processing the received node
1545 // IDs up to the limit. If the request is legitimate then at least they will get
1546 // a response and won't have to resend these nodes in their next request.
1547 tooManyNodeIds = true;
1548 break;
1549 }
1550 auto parsed = deserializeSHAMapNodeID(nodeId);
1551 if (!parsed)
1552 {
1553 peer->charge(resource::kFeeInvalidData, "TMGetLedger: Invalid node ID");
1554 return;
1555 }
1556 nodeIDs.push_back(std::move(*parsed));
1557 }
1558 }
1559
1560 // These are two distinct infractions and are charged independently: requesting too many
1561 // node IDs is charged even for a relay response, while the base "get ledger request" charge
1562 // below is skipped for relay responses.
1563 if (tooManyNodeIds)
1564 {
1565 peer->charge(resource::kFeeModerateBurdenPeer, "TMGetLedger: too many node IDs");
1566
1567 // Truncate the request to what was actually parsed and charged for, so that if this
1568 // request ends up being relayed to another peer, we don't forward the oversized list.
1569 m->mutable_nodeids()->DeleteSubrange(
1570 static_cast<int>(nodeIDs.size()),
1571 m->nodeids_size() - static_cast<int>(nodeIDs.size()));
1572 }
1573 if (!m->has_requestcookie())
1574 {
1575 peer->charge(resource::kFeeModerateBurdenPeer, "TMGetLedger: get ledger request");
1576 }
1577
1578 peer->processLedgerRequest(m, std::move(nodeIDs));
1579 });
1580}
1581
1582void
1584{
1585 JLOG(pJournal_.trace()) << "onMessage, TMProofPathRequest";
1587 {
1588 fee_.update(resource::kFeeMalformedRequest, "proof_path_request disabled");
1589 return;
1590 }
1591
1592 fee_.update(resource::kFeeModerateBurdenPeer, "received a proof path request");
1594 app_.getJobQueue().addJob(JtReplayReq, "RcvProofPReq", [weak, m]() {
1595 if (auto peer = weak.lock())
1596 {
1597 auto reply = peer->ledgerReplayMsgHandler_.processProofPathRequest(m);
1598 if (reply.has_error())
1599 {
1600 if (reply.error() == protocol::TMReplyError::reBAD_REQUEST)
1601 {
1602 peer->charge(resource::kFeeMalformedRequest, "proof_path_request");
1603 }
1604 else
1605 {
1606 peer->charge(resource::kFeeRequestNoReply, "proof_path_request");
1607 }
1608 }
1609 else
1610 {
1611 peer->send(std::make_shared<Message>(reply, protocol::mtPROOF_PATH_RESPONSE));
1612 }
1613 }
1614 });
1615}
1616
1617void
1619{
1621 {
1622 fee_.update(resource::kFeeMalformedRequest, "proof_path_response disabled");
1623 return;
1624 }
1625
1626 switch (ledgerReplayMsgHandler_.processProofPathResponse(m))
1627 {
1629 break;
1631 fee_.update(resource::kFeeInvalidData, "proof_path_response");
1632 break;
1634 fee_.update(resource::kFeeMalformedData, "proof_path_response malformed");
1635 break;
1636 }
1637}
1638
1639void
1641{
1642 JLOG(pJournal_.trace()) << "onMessage, TMReplayDeltaRequest";
1644 {
1645 fee_.update(resource::kFeeMalformedRequest, "replay_delta_request disabled");
1646 return;
1647 }
1648
1651 app_.getJobQueue().addJob(JtReplayReq, "RcvReplDReq", [weak, m]() {
1652 if (auto peer = weak.lock())
1653 {
1654 auto reply = peer->ledgerReplayMsgHandler_.processReplayDeltaRequest(m);
1655 if (reply.has_error())
1656 {
1657 if (reply.error() == protocol::TMReplyError::reBAD_REQUEST)
1658 {
1659 peer->charge(resource::kFeeMalformedRequest, "replay_delta_request");
1660 }
1661 else
1662 {
1663 peer->charge(resource::kFeeRequestNoReply, "replay_delta_request");
1664 }
1665 }
1666 else
1667 {
1668 peer->send(std::make_shared<Message>(reply, protocol::mtREPLAY_DELTA_RESPONSE));
1669 }
1670 }
1671 });
1672}
1673
1674void
1676{
1678 {
1679 fee_.update(resource::kFeeMalformedRequest, "replay_delta_response disabled");
1680 return;
1681 }
1682
1683 switch (ledgerReplayMsgHandler_.processReplayDeltaResponse(m))
1684 {
1686 break;
1688 fee_.update(resource::kFeeInvalidData, "replay_delta_response");
1689 break;
1691 fee_.update(resource::kFeeMalformedData, "replay_delta_response malformed");
1692 break;
1693 }
1694}
1695
1696void
1698{
1699 auto badData = [&](std::string const& msg) {
1700 fee_.update(resource::kFeeInvalidData, msg);
1701 JLOG(pJournal_.warn()) << "TMLedgerData: " << msg;
1702 };
1703
1704 // Verify ledger hash
1705 if (!stringIsUInt256Sized(m->ledgerhash()))
1706 {
1707 badData("Invalid ledger hash");
1708 return;
1709 }
1710
1711 // Verify ledger sequence
1712 {
1713 auto const ledgerSeq{m->ledgerseq()};
1714 if (m->type() == protocol::liTS_CANDIDATE)
1715 {
1716 if (ledgerSeq != 0)
1717 {
1718 badData("Invalid ledger sequence " + std::to_string(ledgerSeq));
1719 return;
1720 }
1721 }
1722 else
1723 {
1724 // Check if within a reasonable range
1725 using namespace std::chrono_literals;
1726 if (app_.getLedgerMaster().getValidatedLedgerAge() <= 10s &&
1727 ledgerSeq > app_.getLedgerMaster().getValidLedgerIndex() + 10)
1728 {
1729 badData("Invalid ledger sequence " + std::to_string(ledgerSeq));
1730 return;
1731 }
1732 }
1733 }
1734
1735 // Verify ledger info type
1736 if (m->type() < protocol::liBASE || m->type() > protocol::liTS_CANDIDATE)
1737 {
1738 badData("Invalid ledger info type");
1739 return;
1740 }
1741
1742 // Verify reply error
1743 if (m->has_error() &&
1744 (m->error() < protocol::reNO_LEDGER || m->error() > protocol::reBAD_REQUEST))
1745 {
1746 badData("Invalid reply error");
1747 return;
1748 }
1749
1750 // Verify ledger nodes.
1751 if (m->nodes_size() <= 0 || m->nodes_size() > tuning::kHardMaxReplyNodes)
1752 {
1753 badData("Invalid Ledger/TXset nodes " + std::to_string(m->nodes_size()));
1754 return;
1755 }
1756
1757 // If there is a request cookie, attempt to relay the message.
1758 if (m->has_requestcookie())
1759 {
1760 if (auto peer = overlay_.findPeerByShortID(m->requestcookie()))
1761 {
1762 m->clear_requestcookie();
1763
1764 // If the original requester doesn't support the new depth-based format, rewrite any
1765 // nodes that use it back to the legacy nodeid format before relaying. Once all nodes
1766 // have upgraded, the old protocol version and this code can be removed. Make sure that
1767 // the format of the nodes is consistent - either all use the legacy format or the new
1768 // format, unless it is liBASE data in which case none of these fields should be set.
1769 auto const peerSupportsNodeDepth =
1770 peer->supportsFeature(ProtocolFeature::LedgerNodeDepth);
1771 enum class MessageType { Unknown, Base, Legacy, Depth };
1772 MessageType messageType = MessageType::Unknown;
1773 for (int i = 0; i < m->nodes_size(); ++i)
1774 {
1775 auto* ledgerNode = m->mutable_nodes(i);
1776
1777 // All nodes should have non-empty data. The field is required so we don't need to
1778 // check for presence first.
1779 if (ledgerNode->nodedata().empty())
1780 {
1781 badData(
1782 "Received node with empty data while relaying ledger data for " +
1783 to_string(UInt256::fromRaw(m->ledgerhash())) + " to peer " +
1784 std::to_string(peer->id()));
1785 return;
1786 }
1787
1788 MessageType msgType = MessageType::Unknown;
1789 if (m->type() == protocol::liBASE)
1790 {
1791 if (ledgerNode->has_nodeid() || ledgerNode->has_id() || ledgerNode->has_depth())
1792 {
1793 badData(
1794 "Received liBASE message with node reference while relaying ledger "
1795 "data for " +
1796 to_string(UInt256::fromRaw(m->ledgerhash())) + " to peer " +
1797 std::to_string(peer->id()));
1798 return;
1799 }
1800 msgType = MessageType::Base;
1801 }
1802 else
1803 {
1804 msgType = ledgerNode->has_nodeid() ? MessageType::Legacy : MessageType::Depth;
1805 }
1806 if (messageType != MessageType::Unknown && messageType != msgType)
1807 {
1808 badData(
1809 "Received mixed mode message while relaying ledger data for " +
1810 to_string(UInt256::fromRaw(m->ledgerhash())) + " to peer " +
1811 std::to_string(peer->id()));
1812 return;
1813 }
1814 messageType = msgType;
1815
1816 if (peerSupportsNodeDepth || msgType != MessageType::Depth)
1817 continue;
1818
1819 SOMETIMES(
1820 !peerSupportsNodeDepth,
1821 "xrpl::PeerImp : relaying depth-format ledger data to pre-2.3 peer");
1822 switch (ledgerNode->reference_case())
1823 {
1824 case protocol::TMLedgerNode::kId: {
1825 // We can directly copy the `id` field, because it uses the same wire format
1826 // as the legacy `nodeid` field.
1827 REACHABLE("xrpl::PeerImp : relay downgrade id to nodeid");
1828 ledgerNode->set_nodeid(ledgerNode->id());
1829 ledgerNode->clear_id();
1830 break;
1831 }
1832 case protocol::TMLedgerNode::kDepth: {
1833 // We need to regenerate the node ID from the node data and depth.
1834 auto treeNode = getTreeNode(ledgerNode->nodedata());
1835 if (!treeNode)
1836 {
1837 badData(
1838 "Unable to get tree node while relaying ledger data for " +
1839 to_string(UInt256::fromRaw(m->ledgerhash())) + " to peer " +
1840 std::to_string(peer->id()));
1841 return;
1842 }
1843
1844 auto const nodeID = getSHAMapNodeID(*ledgerNode, *treeNode);
1845 if (!nodeID)
1846 {
1847 badData(
1848 "Unable to get node ID while relaying ledger data for " +
1849 to_string(UInt256::fromRaw(m->ledgerhash())) + " to peer " +
1850 std::to_string(peer->id()));
1851 return;
1852 }
1853
1854 REACHABLE("xrpl::PeerImp : relay downgrade depth to nodeid");
1855 ledgerNode->set_nodeid(nodeID->getRawString());
1856 ledgerNode->clear_depth();
1857 break;
1858 }
1859 default: {
1860 SOMETIMES(true, "xrpl::PeerImp : relay node has empty reference");
1861 badData(
1862 "Empty node reference while relaying ledger data for " +
1863 to_string(UInt256::fromRaw(m->ledgerhash())) + " to peer " +
1864 std::to_string(peer->id()));
1865 return;
1866 }
1867 }
1868 }
1869
1870 peer->send(std::make_shared<Message>(*m, protocol::mtLEDGER_DATA));
1871 }
1872 else
1873 {
1874 JLOG(pJournal_.info()) << "Unable to route TX/ledger data reply";
1875 }
1876 return;
1877 }
1878
1879 UInt256 const ledgerHash = UInt256::fromRaw(m->ledgerhash());
1880
1881 // Otherwise check if received data for a candidate transaction set
1882 if (m->type() == protocol::liTS_CANDIDATE)
1883 {
1885 app_.getJobQueue().addJob(JtTxnData, "RcvPeerData", [weak, ledgerHash, m]() {
1886 if (auto peer = weak.lock())
1887 {
1888 peer->app_.getInboundTransactions().gotData(ledgerHash, peer, m);
1889 }
1890 });
1891 return;
1892 }
1893
1894 // Consume the message
1895 app_.getInboundLedgers().gotLedgerData(ledgerHash, shared_from_this(), m);
1896}
1897
1898void
1900{
1901 protocol::TMProposeSet const& set = *m;
1902
1903 auto const sig = makeSlice(set.signature());
1904
1905 // Preliminary check for the validity of the signature: A DER encoded
1906 // signature can't be longer than 72 bytes.
1907 if ((std::clamp<std::size_t>(sig.size(), 64, 72) != sig.size()) ||
1908 (publicKeyType(makeSlice(set.nodepubkey())) != KeyType::Secp256k1))
1909 {
1910 JLOG(pJournal_.warn()) << "Proposal: malformed";
1911 fee_.update(resource::kFeeInvalidSignature, " signature can't be longer than 72 bytes");
1912 return;
1913 }
1914
1915 if (!stringIsUInt256Sized(set.currenttxhash()) || !stringIsUInt256Sized(set.previousledger()))
1916 {
1917 JLOG(pJournal_.warn()) << "Proposal: malformed";
1918 fee_.update(resource::kFeeMalformedRequest, "bad hashes");
1919 return;
1920 }
1921
1922 // RH TODO: when isTrusted = false we should probably also cache a key
1923 // suppression for 30 seconds to avoid doing a relatively expensive lookup
1924 // every time a spam packet is received
1925 PublicKey const publicKey{makeSlice(set.nodepubkey())};
1926 auto const isTrusted = app_.getValidators().trusted(publicKey);
1927
1928 // If the operator has specified that untrusted proposals be dropped then
1929 // this happens here I.e. before further wasting CPU verifying the signature
1930 // of an untrusted key
1931 if (!isTrusted)
1932 {
1933 // report untrusted proposal messages
1934 overlay_.reportInboundTraffic(
1936
1937 if (app_.config().relayUntrustedProposals == -1)
1938 return;
1939 }
1940
1941 UInt256 const proposeHash = UInt256::fromRaw(set.currenttxhash());
1942 UInt256 const prevLedger = UInt256::fromRaw(set.previousledger());
1943
1944 NetClock::time_point const closeTime{NetClock::duration{set.closetime()}};
1945
1946 UInt256 const suppression = proposalUniqueId(
1947 proposeHash, prevLedger, set.proposeseq(), closeTime, publicKey.slice(), sig);
1948
1949 if (auto [added, relayed] = app_.getHashRouter().addSuppressionPeerWithStatus(suppression, id_);
1950 !added)
1951 {
1952 // Count unique messages (Slots has it's own 'HashRouter'), which a peer
1953 // receives within IDLED seconds since the message has been relayed.
1954 if (relayed && (stopwatch().now() - *relayed) < reduce_relay::kIdled)
1955 overlay_.updateSlotAndSquelch(suppression, publicKey, id_, protocol::mtPROPOSE_LEDGER);
1956
1957 // report duplicate proposal messages
1958 overlay_.reportInboundTraffic(
1960
1961 JLOG(pJournal_.trace()) << "Proposal: duplicate";
1962
1963 return;
1964 }
1965
1966 if (!isTrusted)
1967 {
1968 if (tracking_.load() == Tracking::Diverged)
1969 {
1970 JLOG(pJournal_.debug()) << "Proposal: Dropping untrusted (peer divergence)";
1971 return;
1972 }
1973
1974 if (!cluster() && app_.getFeeTrack().isLoadedLocal())
1975 {
1976 JLOG(pJournal_.debug()) << "Proposal: Dropping untrusted (load)";
1977 return;
1978 }
1979 }
1980
1981 JLOG(pJournal_.trace()) << "Proposal: " << (isTrusted ? "trusted" : "untrusted");
1982
1983 auto proposal = RCLCxPeerPos(
1984 publicKey,
1985 sig,
1986 suppression,
1988 prevLedger,
1989 set.proposeseq(),
1990 proposeHash,
1991 closeTime,
1992 app_.getTimeKeeper().closeTime(),
1993 calcNodeID(app_.getValidatorManifests().getMasterKey(publicKey))});
1994
1996 app_.getJobQueue().addJob(
1997 isTrusted ? JtProposalT : JtProposalUt, "checkPropose", [weak, isTrusted, m, proposal]() {
1998 if (auto peer = weak.lock())
1999 peer->checkPropose(isTrusted, m, proposal);
2000 });
2001}
2002
2003void
2005{
2006 JLOG(pJournal_.trace()) << "Status: Change";
2007
2008 if (!m->has_networktime())
2009 m->set_networktime(app_.getTimeKeeper().now().time_since_epoch().count());
2010
2011 {
2012 std::scoped_lock const sl(recentLock_);
2013 if (!lastStatus_.has_newstatus() || m->has_newstatus())
2014 {
2015 lastStatus_ = *m;
2016 }
2017 else
2018 {
2019 // preserve old status
2020 protocol::NodeStatus const status = lastStatus_.newstatus();
2021 lastStatus_ = *m;
2022 m->set_newstatus(status);
2023 }
2024 }
2025
2026 if (m->newevent() == protocol::neLOST_SYNC)
2027 {
2028 bool outOfSync{false};
2029 {
2030 // Operations on closedLedgerHash_ and previousLedgerHash_ must be
2031 // guarded by recentLock_.
2032 std::scoped_lock const sl(recentLock_);
2033 if (!closedLedgerHash_.isZero())
2034 {
2035 outOfSync = true;
2036 closedLedgerHash_.zero();
2037 }
2038 previousLedgerHash_.zero();
2039 }
2040 if (outOfSync)
2041 {
2042 JLOG(pJournal_.debug()) << "Status: Out of sync";
2043 }
2044 return;
2045 }
2046
2047 {
2048 UInt256 closedLedgerHash{};
2049 bool const peerChangedLedgers{m->has_ledgerhash() && stringIsUInt256Sized(m->ledgerhash())};
2050
2051 {
2052 // Operations on closedLedgerHash_ and previousLedgerHash_ must be
2053 // guarded by recentLock_.
2054 std::scoped_lock const sl(recentLock_);
2055 if (peerChangedLedgers)
2056 {
2057 closedLedgerHash_ = m->ledgerhash();
2058 closedLedgerHash = closedLedgerHash_;
2059 addLedger(closedLedgerHash, sl);
2060 }
2061 else
2062 {
2063 closedLedgerHash_.zero();
2064 }
2065
2066 if (m->has_ledgerhashprevious() && stringIsUInt256Sized(m->ledgerhashprevious()))
2067 {
2068 previousLedgerHash_ = m->ledgerhashprevious();
2070 }
2071 else
2072 {
2073 previousLedgerHash_.zero();
2074 }
2075 }
2076 if (peerChangedLedgers)
2077 {
2078 JLOG(pJournal_.debug()) << "LCL is " << closedLedgerHash;
2079 }
2080 else
2081 {
2082 JLOG(pJournal_.debug()) << "Status: No ledger";
2083 }
2084 }
2085
2086 if (m->has_firstseq() && m->has_lastseq())
2087 {
2088 std::scoped_lock const sl(recentLock_);
2089
2090 minLedger_ = m->firstseq();
2091 maxLedger_ = m->lastseq();
2092
2093 if ((maxLedger_ < minLedger_) || (minLedger_ == 0) || (maxLedger_ == 0))
2094 minLedger_ = maxLedger_ = 0;
2095 }
2096
2097 if (m->has_ledgerseq() && app_.getLedgerMaster().getValidatedLedgerAge() < 2min)
2098 {
2099 checkTracking(m->ledgerseq(), app_.getLedgerMaster().getValidLedgerIndex());
2100 }
2101
2102 app_.getOPs().pubPeerStatus([m, this]() -> json::Value {
2104
2105 if (m->has_newstatus())
2106 {
2107 switch (m->newstatus())
2108 {
2109 case protocol::nsCONNECTING:
2110 j[jss::status] = "CONNECTING";
2111 break;
2112 case protocol::nsCONNECTED:
2113 j[jss::status] = "CONNECTED";
2114 break;
2115 case protocol::nsMONITORING:
2116 j[jss::status] = "MONITORING";
2117 break;
2118 case protocol::nsVALIDATING:
2119 j[jss::status] = "VALIDATING";
2120 break;
2121 case protocol::nsSHUTTING:
2122 j[jss::status] = "SHUTTING";
2123 break;
2124 }
2125 }
2126
2127 if (m->has_newevent())
2128 {
2129 switch (m->newevent())
2130 {
2131 case protocol::neCLOSING_LEDGER:
2132 j[jss::action] = "CLOSING_LEDGER";
2133 break;
2134 case protocol::neACCEPTED_LEDGER:
2135 j[jss::action] = "ACCEPTED_LEDGER";
2136 break;
2137 case protocol::neSWITCHED_LEDGER:
2138 j[jss::action] = "SWITCHED_LEDGER";
2139 break;
2140 case protocol::neLOST_SYNC:
2141 j[jss::action] = "LOST_SYNC";
2142 break;
2143 }
2144 }
2145
2146 if (m->has_ledgerseq())
2147 {
2148 j[jss::ledger_index] = m->ledgerseq();
2149 }
2150
2151 if (m->has_ledgerhash())
2152 {
2153 UInt256 closedLedgerHash{};
2154 {
2155 std::scoped_lock const sl(recentLock_);
2156 closedLedgerHash = closedLedgerHash_;
2157 }
2158 j[jss::ledger_hash] = to_string(closedLedgerHash);
2159 }
2160
2161 if (m->has_networktime())
2162 {
2163 j[jss::date] = json::UInt(m->networktime());
2164 }
2165
2166 if (m->has_firstseq() && m->has_lastseq())
2167 {
2168 j[jss::ledger_index_min] = json::UInt(m->firstseq());
2169 j[jss::ledger_index_max] = json::UInt(m->lastseq());
2170 }
2171
2172 return j;
2173 });
2174}
2175
2176void
2178{
2179 std::uint32_t serverSeq = 0;
2180 {
2181 // Extract the sequence number of the highest
2182 // ledger this peer has
2183 std::scoped_lock const sl(recentLock_);
2184
2185 serverSeq = maxLedger_;
2186 }
2187 if (serverSeq != 0)
2188 {
2189 // Compare the peer's ledger sequence to the
2190 // sequence of a recently-validated ledger
2191 checkTracking(serverSeq, validationSeq);
2192 }
2193}
2194
2195void
2197{
2198 std::uint32_t const diff = std::max(seq1, seq2) - std::min(seq1, seq2);
2199
2201 {
2202 // The peer's ledger sequence is close to the validation's
2204 }
2205
2206 if ((diff > tuning::kDivergedLedgerLimit) && (tracking_.load() != Tracking::Diverged))
2207 {
2208 // The peer's ledger sequence is way off the validation's
2209 std::scoped_lock const sl(recentLock_);
2210
2213 }
2214}
2215
2216void
2218{
2219 if (!stringIsUInt256Sized(m->hash()))
2220 {
2221 fee_.update(resource::kFeeMalformedRequest, "bad hash");
2222 return;
2223 }
2224
2225 UInt256 const hash = UInt256::fromRaw(m->hash());
2226
2227 if (m->status() == protocol::tsHAVE)
2228 {
2229 std::scoped_lock const sl(recentLock_);
2230
2231 if (std::ranges::find(recentTxSets_, hash) != recentTxSets_.end())
2232 {
2233 fee_.update(resource::kFeeUselessData, "duplicate (tsHAVE)");
2234 return;
2235 }
2236
2237 recentTxSets_.push_back(hash);
2238 }
2239}
2240
2241void
2243 std::string const& messageType,
2244 std::string const& manifest,
2245 std::uint32_t version,
2246 std::vector<ValidatorBlobInfo> const& blobs)
2247{
2248 // If there are no blobs, the message is malformed (possibly because of
2249 // ValidatorList class rules), so charge accordingly and skip processing.
2250 if (blobs.empty())
2251 {
2252 JLOG(pJournal_.warn()) << "Ignored malformed " << messageType;
2253 // This shouldn't ever happen with a well-behaved peer
2254 fee_.update(resource::kFeeHeavyBurdenPeer, "no blobs");
2255 return;
2256 }
2257
2258 auto const hash = sha512Half(manifest, blobs, version);
2259
2260 JLOG(pJournal_.debug()) << "Received " << messageType;
2261
2262 if (!app_.getHashRouter().addSuppressionPeer(hash, id_))
2263 {
2264 JLOG(pJournal_.debug()) << messageType << ": received duplicate " << messageType;
2265 // Charging this fee here won't hurt the peer in the normal
2266 // course of operation (ie. refresh every 5 minutes), but
2267 // will add up if the peer is misbehaving.
2268 fee_.update(resource::kFeeUselessData, "duplicate");
2269 return;
2270 }
2271
2272 auto const applyResult = app_.getValidators().applyListsAndBroadcast(
2273 manifest,
2274 version,
2275 blobs,
2276 remoteAddress_.toString(),
2277 hash,
2278 app_.getOverlay(),
2279 app_.getHashRouter(),
2280 app_.getOPs());
2281
2282 JLOG(pJournal_.debug()) << "Processed " << messageType << " version " << version << " from "
2283 << (applyResult.publisherKey ? strHex(*applyResult.publisherKey)
2284 : "unknown or invalid publisher")
2285 << " with best result " << to_string(applyResult.bestDisposition());
2286
2287 // Act based on the best result
2288 switch (applyResult.bestDisposition())
2289 {
2290 // New list
2292 // Newest list is expired, and that needs to be broadcast, too
2294 // Future list
2297
2298 XRPL_ASSERT(
2299 applyResult.publisherKey,
2300 "xrpl::PeerImp::onValidatorListMessage : publisher key is "
2301 "set");
2302 // NOLINTNEXTLINE(bugprone-unchecked-optional-access) assert above
2303 auto const& pubKey = *applyResult.publisherKey;
2304#ifndef NDEBUG
2305 if (auto const iter = publisherListSequences_.find(pubKey);
2306 iter != publisherListSequences_.end())
2307 {
2308 XRPL_ASSERT(
2309 iter->second < applyResult.sequence,
2310 "xrpl::PeerImp::onValidatorListMessage : lower sequence");
2311 }
2312#endif
2313 publisherListSequences_[pubKey] = applyResult.sequence;
2314 }
2315 break;
2316 // NOLINTNEXTLINE(bugprone-branch-clone): identical to the next branch only in Release
2319#ifndef NDEBUG
2320 {
2322 XRPL_ASSERT(
2323 applyResult.sequence && applyResult.publisherKey,
2324 "xrpl::PeerImp::onValidatorListMessage : nonzero sequence "
2325 "and set publisher key");
2326 XRPL_ASSERT(
2327 publisherListSequences_[*applyResult.publisherKey] <= applyResult.sequence,
2328 "xrpl::PeerImp::onValidatorListMessage : maximum sequence");
2329 }
2330#endif // !NDEBUG
2331
2332 break;
2337 break;
2338 // LCOV_EXCL_START
2339 default:
2340 UNREACHABLE(
2341 "xrpl::PeerImp::onValidatorListMessage : invalid best list "
2342 "disposition");
2343 // LCOV_EXCL_STOP
2344 }
2345
2346 // Charge based on the worst result
2347 switch (applyResult.worstDisposition())
2348 {
2352 // No charges for good data
2353 break;
2356 // Charging this fee here won't hurt the peer in the normal
2357 // course of operation (ie. refresh every 5 minutes), but
2358 // will add up if the peer is misbehaving.
2359 fee_.update(resource::kFeeUselessData, " duplicate (same_sequence or known_sequence)");
2360 break;
2362 // There are very few good reasons for a peer to send an
2363 // old list, particularly more than once.
2364 fee_.update(resource::kFeeInvalidData, "expired");
2365 break;
2367 // Charging this fee here won't hurt the peer in the normal
2368 // course of operation (ie. refresh every 5 minutes), but
2369 // will add up if the peer is misbehaving.
2370 fee_.update(resource::kFeeUselessData, "untrusted");
2371 break;
2373 // This shouldn't ever happen with a well-behaved peer
2374 fee_.update(resource::kFeeInvalidSignature, "invalid list disposition");
2375 break;
2377 // During a version transition, this may be legitimate.
2378 // If it happens frequently, that's probably bad.
2379 fee_.update(resource::kFeeInvalidData, "version");
2380 break;
2381 // LCOV_EXCL_START
2382 default:
2383 UNREACHABLE(
2384 "xrpl::PeerImp::onValidatorListMessage : invalid worst list "
2385 "disposition");
2386 // LCOV_EXCL_STOP
2387 }
2388
2389 // Log based on all the results.
2390 for (auto const& [disp, count] : applyResult.dispositions)
2391 {
2392 switch (disp)
2393 {
2394 // New list
2396 JLOG(pJournal_.debug()) << "Applied " << count << " new " << messageType;
2397 break;
2398 // Newest list is expired, and that needs to be broadcast, too
2400 JLOG(pJournal_.debug()) << "Applied " << count << " expired " << messageType;
2401 break;
2402 // Future list
2404 JLOG(pJournal_.debug()) << "Processed " << count << " future " << messageType;
2405 break;
2407 JLOG(pJournal_.warn())
2408 << "Ignored " << count << " " << messageType << "(s) with current sequence";
2409 break;
2411 JLOG(pJournal_.warn())
2412 << "Ignored " << count << " " << messageType << "(s) with future sequence";
2413 break;
2415 JLOG(pJournal_.warn()) << "Ignored " << count << "stale " << messageType;
2416 break;
2418 JLOG(pJournal_.warn()) << "Ignored " << count << " untrusted " << messageType;
2419 break;
2421 JLOG(pJournal_.warn())
2422 << "Ignored " << count << "unsupported version " << messageType;
2423 break;
2425 JLOG(pJournal_.warn()) << "Ignored " << count << "invalid " << messageType;
2426 break;
2427 // LCOV_EXCL_START
2428 default:
2429 UNREACHABLE(
2430 "xrpl::PeerImp::onValidatorListMessage : invalid list "
2431 "disposition");
2432 // LCOV_EXCL_STOP
2433 }
2434 }
2435}
2436
2437void
2439{
2440 try
2441 {
2442 if (m->version() < 2)
2443 {
2444 JLOG(pJournal_.debug())
2445 << "ValidatorListCollection: received invalid validator list "
2446 "version "
2447 << m->version() << " from peer using protocol version " << to_string(protocol_);
2448 fee_.update(resource::kFeeInvalidData, "wrong version");
2449 return;
2450 }
2452 "ValidatorListCollection", m->manifest(), m->version(), ValidatorList::parseBlobs(*m));
2453 }
2454 catch (std::exception const& e)
2455 {
2456 JLOG(pJournal_.warn()) << "ValidatorListCollection: Exception, " << e.what();
2457 using namespace std::string_literals;
2458 fee_.update(resource::kFeeInvalidData, e.what());
2459 }
2460}
2461
2462void
2464{
2465 if (m->validation().size() < 50)
2466 {
2467 JLOG(pJournal_.warn()) << "Validation: Too small";
2468 fee_.update(resource::kFeeMalformedRequest, "too small");
2469 return;
2470 }
2471
2472 try
2473 {
2474 auto const closeTime = app_.getTimeKeeper().closeTime();
2475
2477 {
2478 SerialIter sit(makeSlice(m->validation()));
2479 try
2480 {
2482 std::ref(sit),
2483 [this](PublicKey const& pk) {
2484 return calcNodeID(app_.getValidatorManifests().getMasterKey(pk));
2485 },
2487 .checkSignature = false, .requireCanonicalOrder = true});
2488 }
2489 catch (std::exception const& e)
2490 {
2491 JLOG(pJournal_.warn()) << "Validation: Exception, " << e.what();
2492 fee_.update(resource::kFeeInvalidData, e.what());
2493 return;
2494 }
2495 val->setSeen(closeTime);
2496 }
2497
2498 if (!isCurrent(
2499 app_.getValidations().parms(),
2500 app_.getTimeKeeper().closeTime(),
2501 val->getSignTime(),
2502 val->getSeenTime()))
2503 {
2504 JLOG(pJournal_.trace()) << "Validation: Not current";
2505 fee_.update(resource::kFeeUselessData, "not current");
2506 return;
2507 }
2508
2509 // RH TODO: when isTrusted = false we should probably also cache a key
2510 // suppression for 30 seconds to avoid doing a relatively expensive
2511 // lookup every time a spam packet is received
2512 auto const isTrusted = app_.getValidators().trusted(val->getSignerPublic());
2513
2514 // If the operator has specified that untrusted validations be
2515 // dropped then this happens here I.e. before further wasting CPU
2516 // verifying the signature of an untrusted key
2517 if (!isTrusted)
2518 {
2519 // increase untrusted validations received
2520 overlay_.reportInboundTraffic(
2522
2523 if (app_.config().relayUntrustedValidations == -1)
2524 return;
2525 }
2526
2527 auto key = sha512Half(makeSlice(m->validation()));
2528
2529 auto [added, relayed] = app_.getHashRouter().addSuppressionPeerWithStatus(key, id_);
2530
2531 if (!added)
2532 {
2533 // Count unique messages (Slots has it's own 'HashRouter'), which a
2534 // peer receives within IDLED seconds since the message has been
2535 // relayed.
2536 if (relayed && (stopwatch().now() - *relayed) < reduce_relay::kIdled)
2537 {
2538 overlay_.updateSlotAndSquelch(
2539 key, val->getSignerPublic(), id_, protocol::mtVALIDATION);
2540 }
2541
2542 // increase duplicate validations received
2543 overlay_.reportInboundTraffic(
2545
2546 JLOG(pJournal_.trace()) << "Validation: duplicate";
2547 return;
2548 }
2549
2550 if (!isTrusted && (tracking_.load() == Tracking::Diverged))
2551 {
2552 JLOG(pJournal_.debug()) << "Dropping untrusted validation from diverged peer";
2553 }
2554 else if (isTrusted || !app_.getFeeTrack().isLoadedLocal())
2555 {
2556 std::string const name = isTrusted ? "ChkTrust" : "ChkUntrust";
2557
2559 app_.getJobQueue().addJob(
2560 isTrusted ? JtValidationT : JtValidationUt, name, [weak, val, m, key]() {
2561 if (auto peer = weak.lock())
2562 peer->checkValidation(val, key, m);
2563 });
2564 }
2565 else
2566 {
2567 JLOG(pJournal_.debug()) << "Dropping untrusted validation for load";
2568 }
2569 }
2570 catch (std::exception const& e)
2571 {
2572 JLOG(pJournal_.warn()) << "Exception processing validation: " << e.what();
2573 using namespace std::string_literals;
2575 }
2576}
2577
2578void
2580{
2581 protocol::TMGetObjectByHash const& packet = *m;
2582
2583 JLOG(pJournal_.trace()) << "received TMGetObjectByHash " << packet.type() << " "
2584 << packet.objects_size();
2585
2586 if (packet.query())
2587 {
2588 // this is a query
2589 if (sendQueue_.size() >= tuning::kDropSendQueue)
2590 {
2591 JLOG(pJournal_.debug()) << "GetObject: Large send queue";
2592 return;
2593 }
2594
2595 if (packet.type() == protocol::TMGetObjectByHash::otFETCH_PACK)
2596 {
2597 doFetchPack(m);
2598 return;
2599 }
2600
2601 if (packet.type() == protocol::TMGetObjectByHash::otTRANSACTIONS)
2602 {
2603 if (!txReduceRelayEnabled())
2604 {
2605 JLOG(pJournal_.error()) << "TMGetObjectByHash: tx reduce-relay is disabled";
2606 fee_.update(resource::kFeeMalformedRequest, "disabled");
2607 return;
2608 }
2609
2611 app_.getJobQueue().addJob(JtRequestedTxn, "DoTxs", [weak, m]() {
2612 if (auto peer = weak.lock())
2613 peer->doTransactions(m);
2614 });
2615 return;
2616 }
2617
2618 if (packet.has_ledgerhash())
2619 {
2620 if (!stringIsUInt256Sized(packet.ledgerhash()))
2621 {
2622 JLOG(pJournal_.debug()) << "GetObj: malformed ledgerhash from peer " << id_;
2623 fee_.update(resource::kFeeMalformedRequest, "get object ledger hash");
2624 return;
2625 }
2626 }
2627 // Reject oversized requests before touching the NodeStore.
2628 // The legitimate upper bound (InboundLedger::getNeededHashes())
2629 // is 8 hashes; anything beyond kHardMaxReplyNodes is non-conforming.
2630 if (packet.objects_size() > tuning::kHardMaxReplyNodes)
2631 {
2632 JLOG(pJournal_.warn())
2633 << "GetObj: oversized request from peer " << id_ << " (" << packet.objects_size()
2634 << " > " << tuning::kHardMaxReplyNodes << ")";
2635 fee_.update(resource::kFeeInvalidData, "oversized get object request");
2636 return;
2637 }
2638
2639 // Dispatch heavy synchronous NodeStore lookups off the peer's
2640 // I/O strand and onto the bounded job queue, mirroring the pattern
2641 // used by processLedgerRequest.
2643 bool const queued = app_.getJobQueue().addJob(JtLedgerReq, "RcvGetObjByHash", [weak, m]() {
2644 auto peer = weak.lock();
2645 if (!peer)
2646 return;
2647 try
2648 {
2649 peer->processGetObjectByHash(m);
2650 }
2651 catch (std::exception const& e)
2652 {
2653 // Surface backend failures (NodeStore I/O, allocation)
2654 // back through the resource model so a misbehaving peer
2655 // is still accountable rather than silently dropped.
2656 JLOG(peer->pJournal_.warn()) << "GetObj: handler threw: " << e.what();
2657 peer->charge(resource::kFeeRequestNoReply, "get object handler exception");
2658 }
2659 });
2660 if (!queued)
2661 {
2662 // The JobQueue is no longer accepting new work (typically
2663 // because it is shutting down / has been joined).
2664 JLOG(pJournal_.warn()) << "GetObj: job queue refused request from peer " << id_;
2665 return;
2666 }
2667
2668 // Admission-time charge: a peer that floods enqueues would
2669 // otherwise be billed only the trivial onMessageEnd fee per
2670 // message until the JobQueue catches up, re-creating an
2671 // uncharged DoS window. Charge the base burden up-front (after
2672 // a successful enqueue); the per-lookup differential is added
2673 // in the worker.
2674 fee_.update(resource::kFeeModerateBurdenPeer, "received a get object by hash request");
2675 }
2676 else
2677 {
2678 // this is a reply
2679 std::uint32_t pLSeq = 0;
2680 bool pLDo = true;
2681 bool progress = false;
2682
2683 for (int i = 0; i < packet.objects_size(); ++i)
2684 {
2685 protocol::TMIndexedObject const& obj = packet.objects(i);
2686
2687 if (obj.has_hash() && stringIsUInt256Sized(obj.hash()))
2688 {
2689 if (obj.has_ledgerseq())
2690 {
2691 if (obj.ledgerseq() != pLSeq)
2692 {
2693 if (pLDo && (pLSeq != 0))
2694 {
2695 JLOG(pJournal_.debug()) << "GetObj: Full fetch pack for " << pLSeq;
2696 }
2697 pLSeq = obj.ledgerseq();
2698 pLDo = !app_.getLedgerMaster().haveLedger(pLSeq);
2699
2700 if (!pLDo)
2701 {
2702 JLOG(pJournal_.debug()) << "GetObj: Late fetch pack for " << pLSeq;
2703 }
2704 else
2705 {
2706 progress = true;
2707 }
2708 }
2709 }
2710
2711 if (pLDo)
2712 {
2713 UInt256 const hash = UInt256::fromRaw(obj.hash());
2714
2715 app_.getLedgerMaster().addFetchPack(
2716 hash, std::make_shared<Blob>(obj.data().begin(), obj.data().end()));
2717 }
2718 }
2719 }
2720
2721 if (pLDo && (pLSeq != 0))
2722 {
2723 JLOG(pJournal_.debug()) << "GetObj: Partial fetch pack for " << pLSeq;
2724 }
2725 if (packet.type() == protocol::TMGetObjectByHash::otFETCH_PACK)
2726 app_.getLedgerMaster().gotFetchPack(progress, pLSeq);
2727 }
2728}
2729
2730void
2732{
2733 protocol::TMGetObjectByHash const& packet = *m;
2734
2735 protocol::TMGetObjectByHash reply;
2736 reply.set_query(false);
2737 reply.set_type(packet.type());
2738
2739 if (packet.has_ledgerhash())
2740 {
2741 reply.set_ledgerhash(packet.ledgerhash());
2742 }
2743
2744 // Defense in depth: caller (onMessage) already validates cheap
2745 // structural properties of the request before dispatching here:
2746 // - objects_size() <= kHardMaxReplyNodes (oversize gate)
2747 // - if has_ledgerhash() then ledgerhash is UInt256-sized
2748 // The iteration cap below mirrors the oversize gate so this method
2749 // remains safe if invoked directly by tests or future callers, and
2750 // a peer cannot drive unbounded NodeStore lookups by sending
2751 // non-existent hashes.
2752 int const requested = packet.objects_size();
2753 int const iterLimit = std::min(requested, tuning::kHardMaxReplyNodes);
2754
2755 for (int i = 0; i < iterLimit; ++i)
2756 {
2757 auto const& obj = packet.objects(i);
2758 if (!obj.has_hash() || !stringIsUInt256Sized(obj.hash()))
2759 continue;
2760
2761 UInt256 const hash = UInt256::fromRaw(obj.hash());
2762 // VFALCO TODO Move this someplace more sensible so we don't
2763 // need to inject the NodeStore interfaces.
2764 std::uint32_t const seq{obj.has_ledgerseq() ? obj.ledgerseq() : 0};
2765 auto const nodeObject = app_.getNodeStore().fetchNodeObject(hash, seq);
2766 if (!nodeObject)
2767 continue;
2768
2769 protocol::TMIndexedObject& newObj = *reply.add_objects();
2770 newObj.set_hash(hash.begin(), hash.size());
2771 auto const& data = nodeObject->getData();
2772 newObj.set_data(data.data(), data.size());
2773 if (obj.has_nodeid())
2774 newObj.set_index(obj.nodeid());
2775 if (obj.has_ledgerseq())
2776 newObj.set_ledgerseq(obj.ledgerseq());
2777 }
2778
2779 // Apply work-proportional charge. `charge()` posts the disconnect
2780 // step (if any) back to strand_, so it is safe to call from this
2781 // JobQueue worker thread.
2782 charge(
2783 // We pass `requested` directly here, instead of actual lookups done. Which could be
2784 // std::min(packet.objects_size(), static_cast<int>(tuning::kHardMaxReplyNodes));
2785 // Because we want to charge as per the request size, to discourage large requests.
2786 computeGetObjectByHashFee(requested, reply.objects_size()),
2787 "processed get object by hash request");
2788
2789 JLOG(pJournal_.trace()) << "GetObj: " << reply.objects_size() << " of " << requested;
2790 send(std::make_shared<Message>(reply, protocol::mtGET_OBJECTS));
2791}
2792
2793void
2795{
2796 if (!txReduceRelayEnabled())
2797 {
2798 JLOG(pJournal_.error()) << "TMHaveTransactions: tx reduce-relay is disabled";
2799 fee_.update(resource::kFeeMalformedRequest, "disabled");
2800 return;
2801 }
2802
2804 app_.getJobQueue().addJob(JtMissingTxn, "HandleHaveTxs", [weak, m]() {
2805 if (auto peer = weak.lock())
2806 peer->handleHaveTransactions(m);
2807 });
2808}
2809
2810void
2812{
2813 protocol::TMGetObjectByHash tmBH;
2814 tmBH.set_type(protocol::TMGetObjectByHash_ObjectType_otTRANSACTIONS);
2815 tmBH.set_query(true);
2816
2817 JLOG(pJournal_.trace()) << "received TMHaveTransactions " << m->hashes_size();
2818
2819 for (std::uint32_t i = 0; i < m->hashes_size(); i++)
2820 {
2821 if (!stringIsUInt256Sized(m->hashes(i)))
2822 {
2823 JLOG(pJournal_.error()) << "TMHaveTransactions with invalid hash size";
2824 fee_.update(resource::kFeeMalformedRequest, "hash size");
2825 return;
2826 }
2827
2828 UInt256 hash = UInt256::fromRaw(m->hashes(i));
2829
2830 auto txn = app_.getMasterTransaction().fetchFromCache(hash);
2831
2832 JLOG(pJournal_.trace()) << "checking transaction " << (bool)txn;
2833
2834 if (!txn)
2835 {
2836 JLOG(pJournal_.debug()) << "adding transaction to request";
2837
2838 auto obj = tmBH.add_objects();
2839 obj->set_hash(hash.data(), hash.size());
2840 }
2841 else
2842 {
2843 // Erase only if a peer has seen this tx. If the peer has not
2844 // seen this tx then the tx could not has been queued for this
2845 // peer.
2846 removeTxQueue(hash);
2847 }
2848 }
2849
2850 JLOG(pJournal_.trace()) << "transaction request object is " << tmBH.objects_size();
2851
2852 if (tmBH.objects_size() > 0)
2853 send(std::make_shared<Message>(tmBH, protocol::mtGET_OBJECTS));
2854}
2855
2856void
2858{
2859 if (!txReduceRelayEnabled())
2860 {
2861 JLOG(pJournal_.error()) << "TMTransactions: tx reduce-relay is disabled";
2862 fee_.update(resource::kFeeMalformedRequest, "disabled");
2863 return;
2864 }
2865
2866 if (m->transactions_size() > reduce_relay::kMaxTxQueueSize)
2867 {
2868 JLOG(pJournal_.error()) << "TMTransactions: transaction list too large";
2869 fee_.update(resource::kFeeMalformedRequest, "Transaction list too large");
2870 return;
2871 }
2872
2873 JLOG(pJournal_.trace()) << "received TMTransactions " << m->transactions_size();
2874
2875 overlay_.addTxMetrics(m->transactions_size());
2876
2877 for (std::uint32_t i = 0; i < m->transactions_size(); ++i)
2878 {
2881 m->mutable_transactions(i), [](protocol::TMTransaction*) {}),
2882 false,
2883 true);
2884 }
2885}
2886
2887void
2889{
2890 dispatch(strand_, [self = shared_from_this(), m]() {
2891 if (!m->has_validatorpubkey())
2892 {
2893 self->fee_.update(resource::kFeeInvalidData, "squelch no pubkey");
2894 return;
2895 }
2896 auto validator = m->validatorpubkey();
2897 auto const slice{makeSlice(validator)};
2898 if (!publicKeyType(slice))
2899 {
2900 self->fee_.update(resource::kFeeInvalidData, "squelch bad pubkey");
2901 return;
2902 }
2903 PublicKey const key(slice);
2904
2905 // Ignore the squelch for validator's own messages.
2906 if (key == self->app_.getValidationPublicKey())
2907 {
2908 JLOG(self->pJournal_.debug())
2909 << "onMessage: TMSquelch discarding validator's squelch " << slice;
2910 return;
2911 }
2912
2913 std::uint32_t const duration = m->has_squelchduration() ? m->squelchduration() : 0;
2914 if (!m->squelch())
2915 {
2916 self->squelch_.removeSquelch(key);
2917 }
2918 else if (!self->squelch_.addSquelch(key, std::chrono::seconds{duration}))
2919 {
2920 self->fee_.update(resource::kFeeInvalidData, "squelch duration");
2921 }
2922
2923 JLOG(self->pJournal_.debug())
2924 << "onMessage: TMSquelch " << slice << " " << self->id() << " " << duration;
2925 });
2926}
2927
2928//--------------------------------------------------------------------------
2929
2930void
2931PeerImp::addLedger(UInt256 const& hash, std::scoped_lock<std::mutex> const& lockedRecentLock)
2932{
2933 // lockedRecentLock is passed as a reminder that recentLock_ must be
2934 // locked by the caller.
2935 (void)lockedRecentLock;
2936
2938 return;
2939
2940 recentLedgers_.push_back(hash);
2941}
2942
2943void
2945{
2946 // VFALCO TODO Invert this dependency using an observer and shared state
2947 // object. Don't queue fetch pack jobs if we're under load or we already
2948 // have some queued.
2949 if (app_.getFeeTrack().isLoadedLocal() ||
2950 (app_.getLedgerMaster().getValidatedLedgerAge() > 40s) ||
2951 (app_.getJobQueue().getJobCount(JtPack) > 10))
2952 {
2953 JLOG(pJournal_.info()) << "Too busy to make fetch pack";
2954 return;
2955 }
2956
2957 if (!stringIsUInt256Sized(packet->ledgerhash()))
2958 {
2959 JLOG(pJournal_.warn()) << "FetchPack hash size malformed";
2960 fee_.update(resource::kFeeMalformedRequest, "hash size");
2961 return;
2962 }
2963
2965
2966 UInt256 const hash = UInt256::fromRaw(packet->ledgerhash());
2967
2969 auto elapsed = UptimeClock::now();
2970 auto const pap = &app_;
2971 app_.getJobQueue().addJob(JtPack, "MakeFetchPack", [pap, weak, packet, hash, elapsed]() {
2972 pap->getLedgerMaster().makeFetchPack(weak, packet, hash, elapsed);
2973 });
2974}
2975
2976void
2978{
2979 protocol::TMTransactions reply;
2980
2981 JLOG(pJournal_.trace()) << "received TMGetObjectByHash requesting tx "
2982 << packet->objects_size();
2983
2984 if (packet->objects_size() > reduce_relay::kMaxTxQueueSize)
2985 {
2986 JLOG(pJournal_.error()) << "doTransactions, invalid number of hashes";
2987 fee_.update(resource::kFeeMalformedRequest, "too big");
2988 return;
2989 }
2990
2991 for (std::uint32_t i = 0; i < packet->objects_size(); ++i)
2992 {
2993 auto const& obj = packet->objects(i);
2994
2995 if (!stringIsUInt256Sized(obj.hash()))
2996 {
2997 fee_.update(resource::kFeeMalformedRequest, "hash size");
2998 return;
2999 }
3000
3001 UInt256 hash = UInt256::fromRaw(obj.hash());
3002
3003 auto txn = app_.getMasterTransaction().fetchFromCache(hash);
3004
3005 if (!txn)
3006 {
3007 JLOG(pJournal_.error())
3008 << "doTransactions, transaction not found " << Slice(hash.data(), hash.size());
3009 fee_.update(resource::kFeeMalformedRequest, "tx not found");
3010 return;
3011 }
3012
3013 Serializer s;
3014 auto tx = reply.add_transactions();
3015 auto sttx = txn->getSTransaction();
3016 sttx->add(s);
3017 tx->set_rawtransaction(s.data(), s.size());
3018 tx->set_status(
3019 txn->getStatus() == TransStatus::INCLUDED ? protocol::tsCURRENT : protocol::tsNEW);
3020 tx->set_receivetimestamp(app_.getTimeKeeper().now().time_since_epoch().count());
3021 tx->set_deferred(txn->getSubmitResult().queued);
3022 }
3023
3024 if (reply.transactions_size() > 0)
3025 send(std::make_shared<Message>(reply, protocol::mtTRANSACTIONS));
3026}
3027
3028void
3030 HashRouterFlags flags,
3031 bool checkSignature,
3032 std::shared_ptr<STTx const> const& stx,
3033 bool batch)
3034{
3035 // VFALCO TODO Rewrite to not use exceptions
3036 try
3037 {
3038 // charge strongly for relaying batch txns
3039 // LCOV_EXCL_START
3040 /*
3041 There is no need to check whether the featureBatchV1_1 amendment is
3042 enabled.
3043
3044 * If the `tfInnerBatchTxn` flag is set, and the amendment is
3045 enabled, then it's an invalid transaction because inner batch
3046 transactions should not be relayed.
3047 * If the `tfInnerBatchTxn` flag is set, and the amendment is *not*
3048 enabled, then the transaction is malformed because it's using an
3049 "unknown" flag. There's no need to waste the resources to send it
3050 to the transaction engine.
3051
3052 We don't normally check transaction validity at this level, but
3053 since we _need_ to check it when the amendment is enabled, we may as
3054 well drop it if the flag is set regardless.
3055 */
3056 if (stx->isFlag(tfInnerBatchTxn))
3057 {
3058 JLOG(pJournal_.warn()) << "Ignoring Network relayed Tx containing "
3059 "tfInnerBatchTxn (checkSignature).";
3060 charge(resource::kFeeModerateBurdenPeer, "inner batch txn");
3061 return;
3062 }
3063 // LCOV_EXCL_STOP
3064
3065 // Expired?
3066 if (stx->isFieldPresent(sfLastLedgerSequence) &&
3067 (stx->getFieldU32(sfLastLedgerSequence) < app_.getLedgerMaster().getValidLedgerIndex()))
3068 {
3069 JLOG(pJournal_.info()) << "Marking transaction " << stx->getTransactionID()
3070 << "as BAD because it's expired";
3071 app_.getHashRouter().setFlags(stx->getTransactionID(), HashRouterFlags::BAD);
3072 charge(resource::kFeeUselessData, "expired tx");
3073 return;
3074 }
3075
3076 if (isPseudoTx(*stx))
3077 {
3078 // Don't do anything with pseudo transactions except put them in the
3079 // TransactionMaster cache
3080 std::string reason;
3081 auto tx = std::make_shared<Transaction>(stx, reason, app_);
3082 XRPL_ASSERT(
3083 tx->getStatus() == TransStatus::NEW,
3084 "xrpl::PeerImp::checkTransaction Transaction created "
3085 "correctly");
3086 if (tx->getStatus() == TransStatus::NEW)
3087 {
3088 JLOG(pJournal_.debug()) << "Processing " << (batch ? "batch" : "unsolicited")
3089 << " pseudo-transaction tx " << tx->getID();
3090
3091 app_.getMasterTransaction().canonicalize(&tx);
3092 // Tell the overlay about it, but don't relay it.
3093 auto const toSkip = app_.getHashRouter().shouldRelay(tx->getID());
3094 if (toSkip)
3095 {
3096 JLOG(pJournal_.debug())
3097 << "Passing skipped pseudo pseudo-transaction tx " << tx->getID();
3098 app_.getOverlay().relay(tx->getID(), {}, *toSkip);
3099 }
3100 if (!batch)
3101 {
3102 JLOG(pJournal_.debug()) << "Charging for pseudo-transaction tx " << tx->getID();
3103 charge(resource::kFeeUselessData, "pseudo tx");
3104 }
3105
3106 return;
3107 }
3108 }
3109
3110 if (checkSignature)
3111 {
3112 // Check the signature before handing off to the job queue.
3113 auto const& validatedRules = app_.getLedgerMaster().getValidatedRules();
3114 if (auto [valid, validReason] =
3115 checkValidity(app_.getHashRouter(), *stx, validatedRules);
3117 {
3118 if (!validReason.empty())
3119 {
3120 JLOG(pJournal_.debug()) << "Exception checking transaction: " << validReason;
3121 }
3122
3123 // For a role-signature transaction, only cache BAD once
3124 // fixCleanup3_4_0 is enabled on this node: the SigBad verdict
3125 // then covers the post-fix prefix and cannot flip back.
3126 // Before the amendment activates, checkValidity's own
3127 // era-scoped cache handles the repeat lookups; setting BAD
3128 // would block a correctly new-prefix-signed transaction until
3129 // the router entry ages out. Remove the guard together with
3130 // the amendment.
3131 if (validatedRules.enabled(fixCleanup3_4_0) ||
3132 (!stx->isFieldPresent(sfSponsorSignature) &&
3133 !stx->isFieldPresent(sfCounterpartySignature)))
3134 {
3135 app_.getHashRouter().setFlags(stx->getTransactionID(), HashRouterFlags::BAD);
3136 }
3137 charge(resource::kFeeInvalidSignature, "check transaction signature failure");
3138 return;
3139 }
3140 }
3141 else
3142 {
3143 forceValidity(app_.getHashRouter(), stx->getTransactionID(), Validity::Valid);
3144 }
3145
3146 std::string reason;
3147 auto tx = std::make_shared<Transaction>(stx, reason, app_);
3148
3149 if (tx->getStatus() == TransStatus::INVALID)
3150 {
3151 if (!reason.empty())
3152 {
3153 JLOG(pJournal_.debug()) << "Exception checking transaction: " << reason;
3154 }
3155 app_.getHashRouter().setFlags(stx->getTransactionID(), HashRouterFlags::BAD);
3156 charge(resource::kFeeInvalidSignature, "tx (impossible)");
3157 return;
3158 }
3159
3160 bool const trusted = any(flags & HashRouterFlags::TRUSTED);
3161 app_.getOPs().processTransaction(tx, trusted, false, NetworkOPs::FailHard::No);
3162 }
3163 catch (std::exception const& ex)
3164 {
3165 JLOG(pJournal_.warn()) << "Exception in " << __func__ << ": " << ex.what();
3166 app_.getHashRouter().setFlags(stx->getTransactionID(), HashRouterFlags::BAD);
3167 using namespace std::string_literals;
3168 charge(resource::kFeeInvalidData, "tx "s + ex.what());
3169 }
3170}
3171
3172// Called from our JobQueue
3173void
3175 bool isTrusted,
3177 RCLCxPeerPos peerPos)
3178{
3179 JLOG(pJournal_.trace()) << "Checking " << (isTrusted ? "trusted" : "UNTRUSTED") << " proposal";
3180
3181 XRPL_ASSERT(packet, "xrpl::PeerImp::checkPropose : non-null packet");
3182
3183 if (!cluster() && !peerPos.checkSign())
3184 {
3185 std::string const desc{"Proposal fails sig check"};
3186 JLOG(pJournal_.warn()) << desc;
3188 return;
3189 }
3190
3191 bool relay = false;
3192
3193 if (isTrusted)
3194 {
3195 relay = app_.getOPs().processTrustedProposal(peerPos);
3196 }
3197 else
3198 {
3199 relay = app_.config().relayUntrustedProposals == 1 || cluster();
3200 }
3201
3202 if (relay)
3203 {
3204 // haveMessage contains peers, which are suppressed; i.e. the peers
3205 // are the source of the message, consequently the message should
3206 // not be relayed to these peers. But the message must be counted
3207 // as part of the squelch logic.
3208 auto haveMessage =
3209 app_.getOverlay().relay(*packet, peerPos.suppressionID(), peerPos.publicKey());
3210 if (!haveMessage.empty())
3211 {
3212 overlay_.updateSlotAndSquelch(
3213 peerPos.suppressionID(),
3214 peerPos.publicKey(),
3215 std::move(haveMessage),
3216 protocol::mtPROPOSE_LEDGER);
3217 }
3218 }
3219}
3220
3221void
3224 UInt256 const& key,
3226{
3227 if (!val->isValid())
3228 {
3229 std::string const desc{"Validation forwarded by peer is invalid"};
3230 JLOG(pJournal_.debug()) << desc;
3232 return;
3233 }
3234
3235 // FIXME it should be safe to remove this try/catch. Investigate codepaths.
3236 try
3237 {
3238 if (app_.getOPs().recvValidation(val, std::to_string(id())) || cluster())
3239 {
3240 // haveMessage contains peers, which are suppressed; i.e. the peers
3241 // are the source of the message, consequently the message should
3242 // not be relayed to these peers. But the message must be counted
3243 // as part of the squelch logic.
3244 auto haveMessage = overlay_.relay(*packet, key, val->getSignerPublic());
3245 if (!haveMessage.empty())
3246 {
3247 overlay_.updateSlotAndSquelch(
3248 key, val->getSignerPublic(), std::move(haveMessage), protocol::mtVALIDATION);
3249 }
3250 }
3251 }
3252 catch (std::exception const& ex)
3253 {
3254 JLOG(pJournal_.trace()) << "Exception processing validation: " << ex.what();
3255 using namespace std::string_literals;
3256 charge(resource::kFeeMalformedRequest, "validation "s + ex.what());
3257 }
3258}
3259
3260// Returns the set of peers that can help us get
3261// the TX tree with the specified root hash.
3262//
3264getPeerWithTree(OverlayImpl& ov, UInt256 const& rootHash, PeerImp const* skip)
3265{
3267 int retScore = 0;
3268
3269 ov.forEach([&](std::shared_ptr<PeerImp>&& p) {
3270 if (p->hasTxSet(rootHash) && p.get() != skip)
3271 {
3272 auto score = p->getScore(true);
3273 if (!ret || (score > retScore))
3274 {
3275 ret = std::move(p);
3276 retScore = score;
3277 }
3278 }
3279 });
3280
3281 return ret;
3282}
3283
3284// Returns a random peer weighted by how likely to
3285// have the ledger and how responsive it is.
3286//
3289 OverlayImpl& ov,
3290 UInt256 const& ledgerHash,
3291 LedgerIndex ledger,
3292 PeerImp const* skip)
3293{
3295 int retScore = 0;
3296
3297 ov.forEach([&](std::shared_ptr<PeerImp>&& p) {
3298 if (p->hasLedger(ledgerHash, ledger) && p.get() != skip)
3299 {
3300 auto score = p->getScore(true);
3301 if (!ret || (score > retScore))
3302 {
3303 ret = std::move(p);
3304 retScore = score;
3305 }
3306 }
3307 });
3308
3309 return ret;
3310}
3311
3312void
3314 std::shared_ptr<Ledger const> const& ledger,
3315 protocol::TMLedgerData& ledgerData)
3316{
3317 JLOG(pJournal_.trace()) << "sendLedgerBase: Base data";
3318
3319 Serializer s(sizeof(LedgerHeader));
3320 addRaw(ledger->header(), s);
3321 ledgerData.add_nodes()->set_nodedata(s.getDataPtr(), s.getLength());
3322
3323 auto const& stateMap{ledger->stateMap()};
3324 if (stateMap.getHash() != beast::kZero)
3325 {
3326 // Return account state root node if possible
3327 Serializer root(768);
3328
3329 stateMap.serializeRoot(root);
3330 ledgerData.add_nodes()->set_nodedata(root.getDataPtr(), root.getLength());
3331
3332 if (ledger->header().txHash != beast::kZero)
3333 {
3334 auto const& txMap{ledger->txMap()};
3335 if (txMap.getHash() != beast::kZero)
3336 {
3337 // Return TX root node if possible
3338 root.erase();
3339 txMap.serializeRoot(root);
3340 ledgerData.add_nodes()->set_nodedata(root.getDataPtr(), root.getLength());
3341 }
3342 }
3343 }
3344
3345 auto message{std::make_shared<Message>(ledgerData, protocol::mtLEDGER_DATA)};
3346 send(message);
3347}
3348
3351{
3352 JLOG(pJournal_.trace()) << "getLedger: Ledger";
3353
3355
3356 if (m->has_ledgerhash())
3357 {
3358 // Attempt to find ledger by hash
3359 UInt256 const ledgerHash = UInt256::fromRaw(m->ledgerhash());
3360 ledger = app_.getLedgerMaster().getLedgerByHash(ledgerHash);
3361 if (!ledger)
3362 {
3363 JLOG(pJournal_.trace()) << "getLedger: Don't have ledger with hash " << ledgerHash;
3364
3365 if (m->has_querytype() && !m->has_requestcookie())
3366 {
3367 // Attempt to relay the request to a peer
3368 if (auto const peer = getPeerWithLedger(
3369 overlay_, ledgerHash, m->has_ledgerseq() ? m->ledgerseq() : 0, this))
3370 {
3371 m->set_requestcookie(id());
3372 peer->send(std::make_shared<Message>(*m, protocol::mtGET_LEDGER));
3373 JLOG(pJournal_.debug()) << "getLedger: Request relayed to peer";
3374 return ledger;
3375 }
3376
3377 JLOG(pJournal_.trace()) << "getLedger: Failed to find peer to relay request";
3378 }
3379 }
3380 }
3381 else if (m->has_ledgerseq())
3382 {
3383 // Attempt to find ledger by sequence
3384 if (m->ledgerseq() < app_.getLedgerMaster().getEarliestFetch())
3385 {
3386 JLOG(pJournal_.debug()) << "getLedger: Early ledger sequence request";
3387 }
3388 else
3389 {
3390 ledger = app_.getLedgerMaster().getLedgerBySeq(m->ledgerseq());
3391 if (!ledger)
3392 {
3393 JLOG(pJournal_.debug())
3394 << "getLedger: Don't have ledger with sequence " << m->ledgerseq();
3395 }
3396 }
3397 }
3398 else if (m->has_ltype() && m->ltype() == protocol::ltCLOSED)
3399 {
3400 ledger = app_.getLedgerMaster().getClosedLedger();
3401 }
3402
3403 if (ledger)
3404 {
3405 // Validate retrieved ledger sequence
3406 auto const ledgerSeq{ledger->header().seq};
3407 if (m->has_ledgerseq())
3408 {
3409 if (ledgerSeq != m->ledgerseq())
3410 {
3411 // Do not resource charge a peer responding to a relay
3412 if (!m->has_requestcookie())
3413 charge(resource::kFeeMalformedRequest, "get_ledger ledgerSeq");
3414
3415 ledger.reset();
3416 JLOG(pJournal_.warn()) << "getLedger: Invalid ledger sequence " << ledgerSeq;
3417 }
3418 }
3419 else if (ledgerSeq < app_.getLedgerMaster().getEarliestFetch())
3420 {
3421 ledger.reset();
3422 JLOG(pJournal_.debug()) << "getLedger: Early ledger sequence request " << ledgerSeq;
3423 }
3424 }
3425 else
3426 {
3427 JLOG(pJournal_.debug()) << "getLedger: Unable to find ledger";
3428 }
3429
3430 return ledger;
3431}
3432
3435{
3436 JLOG(pJournal_.trace()) << "getTxSet: TX set";
3437
3438 UInt256 const txSetHash = UInt256::fromRaw(m->ledgerhash());
3439 std::shared_ptr<SHAMap> shaMap{app_.getInboundTransactions().getSet(txSetHash, false)};
3440 if (!shaMap)
3441 {
3442 if (m->has_querytype() && !m->has_requestcookie())
3443 {
3444 // Attempt to relay the request to a peer
3445 if (auto const peer = getPeerWithTree(overlay_, txSetHash, this))
3446 {
3447 m->set_requestcookie(id());
3448 peer->send(std::make_shared<Message>(*m, protocol::mtGET_LEDGER));
3449 JLOG(pJournal_.debug()) << "getTxSet: Request relayed";
3450 }
3451 else
3452 {
3453 JLOG(pJournal_.debug()) << "getTxSet: Failed to find relay peer";
3454 }
3455 }
3456 else
3457 {
3458 JLOG(pJournal_.debug()) << "getTxSet: Failed to find TX set";
3459 }
3460 }
3461
3462 return shaMap;
3463}
3464
3465void
3469{
3472 SHAMap const* map{nullptr};
3473 protocol::TMLedgerData ledgerData;
3474 bool fatLeaves{true};
3475 auto const itype{m->itype()};
3476
3477 if (itype == protocol::liTS_CANDIDATE)
3478 {
3479 if (sharedMap = getTxSet(m); !sharedMap)
3480 return;
3481 map = sharedMap.get();
3482
3483 // Fill out the reply
3484 ledgerData.set_ledgerseq(0);
3485 ledgerData.set_ledgerhash(m->ledgerhash());
3486 ledgerData.set_type(protocol::liTS_CANDIDATE);
3487 if (m->has_requestcookie())
3488 ledgerData.set_requestcookie(m->requestcookie());
3489
3490 // We'll already have most transactions
3491 fatLeaves = false;
3492 }
3493 else
3494 {
3495 if (sendQueue_.size() >= tuning::kDropSendQueue)
3496 {
3497 JLOG(pJournal_.debug()) << "processLedgerRequest: Large send queue";
3498 return;
3499 }
3500 if (app_.getFeeTrack().isLoadedLocal() && !cluster())
3501 {
3502 JLOG(pJournal_.debug()) << "processLedgerRequest: Too busy";
3503 return;
3504 }
3505
3506 if (ledger = getLedger(m); !ledger)
3507 return;
3508
3509 // Fill out the reply
3510 auto const ledgerHash{ledger->header().hash};
3511 ledgerData.set_ledgerhash(ledgerHash.begin(), ledgerHash.size());
3512 ledgerData.set_ledgerseq(ledger->header().seq);
3513 ledgerData.set_type(itype);
3514 if (m->has_requestcookie())
3515 ledgerData.set_requestcookie(m->requestcookie());
3516
3517 switch (itype)
3518 {
3519 case protocol::liBASE:
3520 sendLedgerBase(ledger, ledgerData);
3521 return;
3522
3523 case protocol::liTX_NODE:
3524 map = &ledger->txMap();
3525 JLOG(pJournal_.trace())
3526 << "processLedgerRequest: TX map hash " << to_string(map->getHash());
3527 break;
3528
3529 case protocol::liAS_NODE:
3530 map = &ledger->stateMap();
3531 JLOG(pJournal_.trace())
3532 << "processLedgerRequest: Account state map hash " << to_string(map->getHash());
3533 break;
3534
3535 default:
3536 // This case should not be possible here
3537 JLOG(pJournal_.error()) << "processLedgerRequest: Invalid ledger info type";
3538 return;
3539 }
3540 }
3541
3542 if (map == nullptr)
3543 {
3544 JLOG(pJournal_.warn()) << "processLedgerRequest: Unable to find map";
3545 return;
3546 }
3547
3548 // Add requested node data to reply
3549 if (!nodeIDs.empty())
3550 {
3551 std::uint32_t const defaultDepth = isHighLatency() ? 2 : 1;
3552 auto const queryDepth{m->has_querydepth() ? m->querydepth() : defaultDepth};
3553
3555 data.reserve(tuning::kSoftMaxReplyNodes);
3556 auto const useLedgerNodeDepth = supportsFeature(ProtocolFeature::LedgerNodeDepth);
3557
3558 for (auto const& nodeID : nodeIDs)
3559 {
3560 if (ledgerData.nodes_size() >= tuning::kSoftMaxReplyNodes)
3561 break;
3562
3563 data.clear();
3564
3565 try
3566 {
3567 if (map->getNodeFat(nodeID, data, fatLeaves, queryDepth))
3568 {
3569 JLOG(pJournal_.trace())
3570 << "processLedgerRequest: getNodeFat got " << data.size() << " nodes";
3571
3572 for (auto const& d : data)
3573 {
3574 if (ledgerData.nodes_size() >= tuning::kHardMaxReplyNodes)
3575 break;
3576
3577 protocol::TMLedgerNode* node{ledgerData.add_nodes()};
3578 node->set_nodedata(d.data.data(), d.data.size());
3579
3580 // When the LedgerNodeDepth protocol feature is not supported by the peer,
3581 // we always set the `nodeid` field. However, when it is supported then we
3582 // set the `id` field for inner nodes and the `depth` field for leaf nodes.
3583 if (!useLedgerNodeDepth)
3584 {
3585 node->set_nodeid(d.nodeID.getRawString());
3586 }
3587 else if (d.isLeaf)
3588 {
3589 REACHABLE("xrpl::PeerImp : emit leaf depth in reply");
3590 node->set_depth(d.nodeID.getDepth());
3591 }
3592 else
3593 {
3594 REACHABLE("xrpl::PeerImp : emit inner id in reply");
3595 node->set_id(d.nodeID.getRawString());
3596 }
3597 }
3598 }
3599 else
3600 {
3601 JLOG(pJournal_.warn()) << "processLedgerRequest: getNodeFat returns false";
3602 }
3603 }
3604 catch (std::exception const& e)
3605 {
3606 std::string info;
3607 switch (itype)
3608 {
3609 case protocol::liBASE:
3610 // This case should not be possible here
3611 info = "Ledger base";
3612 break;
3613
3614 case protocol::liTX_NODE:
3615 info = "TX node";
3616 break;
3617
3618 case protocol::liAS_NODE:
3619 info = "AS node";
3620 break;
3621
3622 case protocol::liTS_CANDIDATE:
3623 info = "TS candidate";
3624 break;
3625
3626 default:
3627 info = "Invalid";
3628 break;
3629 }
3630
3631 if (!m->has_ledgerhash())
3632 info += ", no hash specified";
3633
3634 JLOG(pJournal_.warn())
3635 << "processLedgerRequest: getNodeFat with nodeId " << nodeID
3636 << " and ledger info type " << info << " throws exception: " << e.what();
3637 }
3638 }
3639
3640 JLOG(pJournal_.info()) << "processLedgerRequest: Got request for " << m->nodeids_size()
3641 << " node IDs at depth " << queryDepth << ", return "
3642 << ledgerData.nodes_size() << " nodes";
3643 }
3644
3645 if (ledgerData.nodes_size() == 0)
3646 return;
3647
3648 send(std::make_shared<Message>(ledgerData, protocol::mtLEDGER_DATA));
3649}
3650
3651// Differential pricing helper. Returns only the *dynamic* component
3652// of the per-message charge — the base `kFeeModerateBurdenPeer` is
3653// applied at admission time in `onMessage(TMGetObjectByHash)` so a
3654// high traffic client pays for the message regardless of when (or
3655// whether) the worker runs.
3656//
3657// Dynamic charge model:
3658//
3659// billable = max(0, requested - kFreeObjectsPerRequest)
3660// missed = max(0, requested - found)
3661// billableMisses = min(missed, billable) // misses billed first
3662// billableHits = billable - billableMisses
3663// sizeBand = (requested > kBandMediumMax) ? kCostBandLarge
3664// : (requested > kBandSmallMax) ? kCostBandMedium
3665// : kCostBandSmall
3666// dynamic = billableHits * kCostPerLookupHit
3667// + billableMisses * kCostPerLookupMiss
3668// + sizeBand
3669//
3670// Misses are billed first against the billable budget because a node store
3671// seek dominates a cache hit and because invalid hashes are ~100% miss by construction.
3673PeerImp::computeGetObjectByHashFee(int const requested, int const found)
3674{
3675 int const billable = std::max(0, requested - static_cast<int>(tuning::kFreeObjectsPerRequest));
3676 // Clamp `missed` so a future caller passing found > requested cannot
3677 // produce a negative value that flips the hits/misses split.
3678 int const missed = std::max(0, requested - found);
3679 int const billableMisses = std::min(missed, billable);
3680 int const billableHits = billable - billableMisses;
3681
3682 int sizeBand = tuning::kCostBandSmall;
3683 if (requested > tuning::kBandMediumMax)
3684 {
3685 sizeBand = tuning::kCostBandLarge;
3686 }
3687 else if (requested > tuning::kBandSmallMax)
3688 {
3689 sizeBand = tuning::kCostBandMedium;
3690 }
3691
3692 int const dynamic = (billableHits * tuning::kCostPerLookupHit) +
3693 (billableMisses * tuning::kCostPerLookupMiss) + sizeBand;
3694
3695 return resource::Charge(dynamic, "GetObject differential");
3696}
3697
3698int
3699PeerImp::getScore(bool haveItem) const
3700{
3701 // Random component of score, used to break ties and avoid
3702 // overloading the "best" peer
3703 static int const kSpRandomMax = 9999;
3704
3705 // Score for being very likely to have the thing we are
3706 // look for; should be roughly spRandomMax
3707 static int const kSpHaveItem = 10000;
3708
3709 // Score reduction for each millisecond of latency; should
3710 // be roughly spRandomMax divided by the maximum reasonable
3711 // latency
3712 static int const kSpLatency = 30;
3713
3714 // Penalty for unknown latency; should be roughly spRandomMax
3715 static int const kSpNoLatency = 8000;
3716
3717 int score = randInt(kSpRandomMax);
3718
3719 if (haveItem)
3720 score += kSpHaveItem;
3721
3723 {
3724 std::scoped_lock const sl(recentLock_);
3725 latency = latency_;
3726 }
3727
3728 if (latency)
3729 {
3730 score -= latency->count() * kSpLatency;
3731 }
3732 else
3733 {
3734 score -= kSpNoLatency;
3735 }
3736
3737 return score;
3738}
3739
3740bool
3742{
3743 std::scoped_lock const sl(recentLock_);
3744 return latency_ >= kPeerHighLatency;
3745}
3746
3747void
3749{
3750 using namespace std::chrono_literals;
3751 std::unique_lock const lock{mutex_};
3752
3753 totalBytes_ += bytes;
3754 accumBytes_ += bytes;
3755 auto const timeElapsed = ClockType::now() - intervalStart_;
3756 auto const timeElapsedInSecs = std::chrono::duration_cast<std::chrono::seconds>(timeElapsed);
3757
3758 if (timeElapsedInSecs >= 1s)
3759 {
3760 auto const avgBytes = accumBytes_ / timeElapsedInSecs.count();
3761 rollingAvg_.push_back(avgBytes);
3762
3763 auto const totalBytes = std::accumulate(rollingAvg_.begin(), rollingAvg_.end(), 0ull);
3765
3767 accumBytes_ = 0;
3768 }
3769}
3770
3773{
3774 std::shared_lock const lock{mutex_};
3775 return rollingAvgBytes_;
3776}
3777
3780{
3781 std::shared_lock const lock{mutex_};
3782 return totalBytes_;
3783}
3784
3785} // namespace xrpl
T accumulate(T... args)
T begin(T... args)
T clamp(T... args)
A version-independent IP address and port combination.
Definition IPEndpoint.h:24
static std::optional< Endpoint > fromStringChecked(std::string const &s)
Create an Endpoint from a string.
static Endpoint fromString(std::string const &s)
Represents a JSON value.
Definition json_value.h:117
static BaseUInt fromRaw(Container const &c)
Definition base_uint.h:302
pointer data()
Definition base_uint.h:117
iterator begin()
Definition base_uint.h:128
static constexpr std::size_t size()
Definition base_uint.h:548
static std::size_t messageSize(::google::protobuf::Message const &message)
Definition Message.cpp:48
std::chrono::time_point< NetClock > time_point
Definition chrono.h:48
std::chrono::duration< rep, period > duration
Definition chrono.h:47
Child(OverlayImpl &overlay)
void forEach(UnaryFunc &&f) const
std::shared_mutex mutex_
Definition PeerImp.h:237
std::uint64_t rollingAvgBytes_
Definition PeerImp.h:242
boost::circular_buffer< std::uint64_t > rollingAvg_
Definition PeerImp.h:238
void addMessage(std::uint64_t bytes)
Definition PeerImp.cpp:3748
std::uint64_t averageBytes() const
Definition PeerImp.cpp:3772
ClockType::time_point intervalStart_
Definition PeerImp.h:239
std::uint64_t totalBytes() const
Definition PeerImp.cpp:3779
std::uint64_t accumBytes_
Definition PeerImp.h:241
std::uint64_t totalBytes_
Definition PeerImp.h:240
void checkTracking(std::uint32_t validationSeq)
Check if the peer is tracking.
Definition PeerImp.cpp:2177
void onTimer(boost::system::error_code const &ec)
Definition PeerImp.cpp:702
std::optional< std::chrono::milliseconds > latency_
Definition PeerImp.h:129
std::string getVersion() const
Return the version of xrpld that the peer is running, if reported.
Definition PeerImp.cpp:419
void onMessage(std::shared_ptr< protocol::TMManifests > const &m)
Definition PeerImp.cpp:1095
ProtocolVersion protocol_
Definition PeerImp.h:109
bool txReduceRelayEnabled_
Definition PeerImp.h:211
void checkTransaction(HashRouterFlags flags, bool checkSignature, std::shared_ptr< STTx const > const &stx, bool batch)
Definition PeerImp.cpp:3029
void setTimer()
Definition PeerImp.cpp:662
void handleTransaction(std::shared_ptr< protocol::TMTransaction > const &m, bool eraseTxQueue, bool batch)
Called from onMessage(TMTransaction(s)).
Definition PeerImp.cpp:1293
bool hasTxSet(UInt256 const &hash) const override
Definition PeerImp.cpp:579
beast::WrappedSink sink_
Definition PeerImp.h:89
void addLedger(UInt256 const &hash, std::scoped_lock< std::mutex > const &lockedRecentLock)
Definition PeerImp.cpp:2931
std::string name() const
Definition PeerImp.cpp:862
bool txReduceRelayEnabled() const override
Definition PeerImp.h:487
SocketType & socket_
Definition PeerImp.h:94
compression::Compressed Compressed
Definition PeerImp.h:83
boost::beast::http::fields const & headers_
Definition PeerImp.h:195
Compressed compressionEnabled_
Definition PeerImp.h:204
ClockType::duration uptime() const
Definition PeerImp.h:399
StreamType & stream_
Definition PeerImp.h:95
ClockType::time_point const creationTime_
Definition PeerImp.h:132
boost::circular_buffer< UInt256 > recentTxSets_
Definition PeerImp.h:127
LedgerIndex minLedger_
Definition PeerImp.h:121
std::string prefix_
Definition PeerImp.h:88
WaitableTimer timer_
Definition PeerImp.h:97
std::shared_ptr< peer_finder::Slot > const slot_
Definition PeerImp.h:191
void onWriteMessage(ErrorCode ec, std::size_t bytesTransferred)
Definition PeerImp.cpp:989
boost::beast::multi_buffer readBuffer_
Definition PeerImp.h:192
void onShutdown(ErrorCode ec)
Definition PeerImp.cpp:762
boost::circular_buffer< UInt256 > recentLedgers_
Definition PeerImp.h:126
std::string const & fingerprint() const override
Definition PeerImp.h:573
void checkValidation(std::shared_ptr< STValidation > const &val, UInt256 const &key, std::shared_ptr< protocol::TMValidation > const &packet)
Definition PeerImp.cpp:3222
bool ledgerReplayEnabled_
Definition PeerImp.h:213
void sendTxQueue() override
Send aggregated transactions' hashes.
Definition PeerImp.cpp:338
struct xrpl::PeerImp::@337373043150231020277011015352151251117171316327 metrics_
beast::ip::Endpoint const remoteAddress_
Definition PeerImp.h:101
reduce_relay::Squelch< UptimeClock > squelch_
Definition PeerImp.h:134
PeerImp(PeerImp const &)=delete
void cycleStatus() override
Definition PeerImp.cpp:586
bool gracefulClose_
Definition PeerImp.h:197
std::shared_mutex nameMutex_
Definition PeerImp.h:117
std::string domain() const
Definition PeerImp.cpp:869
std::atomic< Tracking > tracking_
Definition PeerImp.h:111
void ledgerRange(std::uint32_t &minSeq, std::uint32_t &maxSeq) const override
Definition PeerImp.cpp:570
HashMap< PublicKey, std::size_t > publisherListSequences_
Definition PeerImp.h:202
void addTxQueue(UInt256 const &hash) override
Add transaction's hash to the transactions' hashes queue.
Definition PeerImp.cpp:354
void onReadMessage(ErrorCode ec, std::size_t bytesTransferred)
Definition PeerImp.cpp:917
LedgerIndex maxLedger_
Definition PeerImp.h:122
beast::WrappedSink pSink_
Definition PeerImp.h:90
Application & app_
Definition PeerImp.h:85
PublicKey const publicKey_
Definition PeerImp.h:115
void processGetObjectByHash(std::shared_ptr< protocol::TMGetObjectByHash > const &m)
Process a generic-query TMGetObjectByHash message.
Definition PeerImp.cpp:2731
bool cluster() const override
Returns true if this connection is a member of the cluster.
Definition PeerImp.cpp:413
static resource::Charge computeGetObjectByHashFee(int const requested, int const found)
Compute the per-message resource charge for a TMGetObjectByHash request based on how much work was ac...
Definition PeerImp.cpp:3673
std::queue< std::shared_ptr< Message > > sendQueue_
Definition PeerImp.h:196
void checkPropose(bool isTrusted, std::shared_ptr< protocol::TMProposeSet > const &packet, RCLCxPeerPos peerPos)
Definition PeerImp.cpp:3174
virtual void run()
Definition PeerImp.cpp:206
OverlayImpl & overlay_
Definition PeerImp.h:105
std::chrono::steady_clock ClockType
Definition PeerImp.h:75
std::unique_ptr< StreamType > streamPtr_
Definition PeerImp.h:93
resource::Consumer usage_
Definition PeerImp.h:183
void sendLedgerBase(std::shared_ptr< Ledger const > const &ledger, protocol::TMLedgerData &ledgerData)
Definition PeerImp.cpp:3313
bool hasLedger(UInt256 const &hash, std::uint32_t seq) const override
Definition PeerImp.cpp:556
beast::Journal const pJournal_
Definition PeerImp.h:92
ClockType::time_point trackingTime_
Definition PeerImp.h:112
friend class OverlayImpl
Definition PeerImp.h:216
static std::string makePrefix(std::string const &fingerprint)
Definition PeerImp.cpp:694
ChargeWithContext fee_
Definition PeerImp.h:184
void handleHaveTransactions(std::shared_ptr< protocol::TMHaveTransactions > const &m)
Handle protocol message with hashes of transactions that have not been relayed by an upstream node do...
Definition PeerImp.cpp:2811
void onMessageUnknown(std::uint16_t type)
Definition PeerImp.cpp:1041
void processLedgerRequest(std::shared_ptr< protocol::TMGetLedger > const &m, std::vector< SHAMapNodeID > nodeIDs)
Definition PeerImp.cpp:3466
~PeerImp() override
Definition PeerImp.cpp:183
void gracefulClose()
Definition PeerImp.cpp:647
Tracking
Whether the peer's view of the ledger converges or diverges from ours.
Definition PeerImp.h:72
protocol::TMStatusChange lastStatus_
Definition PeerImp.h:182
LedgerReplayMsgHandler ledgerReplayMsgHandler_
Definition PeerImp.h:214
boost::asio::basic_waitable_timer< std::chrono::steady_clock > WaitableTimer
Definition PeerImp.h:82
void fail(std::string const &reason)
Definition PeerImp.cpp:622
std::unique_ptr< LoadEvent > loadEvent_
Definition PeerImp.h:199
void send(std::shared_ptr< Message > const &m) override
Definition PeerImp.cpp:277
void onMessageBegin(std::uint16_t type, std::shared_ptr<::google::protobuf::Message > const &m, std::size_t size, std::size_t uncompressedSize, bool isCompressed)
Definition PeerImp.cpp:1047
UInt256 closedLedgerHash_
Definition PeerImp.h:123
void doFetchPack(std::shared_ptr< protocol::TMGetObjectByHash > const &packet)
Definition PeerImp.cpp:2944
bool supportsFeature(ProtocolFeature f) const override
Definition PeerImp.cpp:541
void doProtocolStart()
Definition PeerImp.cpp:879
int largeSendq_
Definition PeerImp.h:198
std::shared_ptr< Ledger const > getLedger(std::shared_ptr< protocol::TMGetLedger > const &m)
Definition PeerImp.cpp:3350
bool const inbound_
Definition PeerImp.h:106
boost::asio::strand< boost::asio::executor > strand_
Definition PeerImp.h:96
int getScore(bool haveItem) const override
Definition PeerImp.cpp:3699
std::string name_
Definition PeerImp.h:116
beast::Journal const journal_
Definition PeerImp.h:91
std::string fingerprint_
Definition PeerImp.h:87
bool isHighLatency() const override
Definition PeerImp.cpp:3741
void cancelTimer() noexcept
Definition PeerImp.cpp:679
ID const id_
Definition PeerImp.h:86
bool hasRange(std::uint32_t uMin, std::uint32_t uMax) override
Definition PeerImp.cpp:596
std::shared_ptr< SHAMap const > getTxSet(std::shared_ptr< protocol::TMGetLedger > const &m) const
Definition PeerImp.cpp:3434
std::optional< std::uint32_t > lastPingSeq_
Definition PeerImp.h:130
bool detaching_
Definition PeerImp.h:113
boost::system::error_code ErrorCode
Definition PeerImp.h:76
bool crawl() const
Returns true if this connection will publicly share its IP address.
Definition PeerImp.cpp:404
void stop() override
Definition PeerImp.cpp:264
HttpRequestType request_
Definition PeerImp.h:193
void onValidatorListMessage(std::string const &messageType, std::string const &manifest, std::uint32_t version, std::vector< ValidatorBlobInfo > const &blobs)
Definition PeerImp.cpp:2242
Peer::ID id() const override
Definition PeerImp.h:361
void charge(resource::Charge const &fee, std::string const &context) override
Adjust this peer's load balance based on the type of load imposed.
Definition PeerImp.cpp:378
void removeTxQueue(UInt256 const &hash) override
Remove transaction's hash from the transactions' hashes queue.
Definition PeerImp.cpp:369
void doTransactions(std::shared_ptr< protocol::TMGetObjectByHash > const &packet)
Process peer's request to send missing transactions.
Definition PeerImp.cpp:2977
void doAccept()
Definition PeerImp.cpp:788
std::shared_ptr< peer_finder::Slot > const & slot()
Definition PeerImp.h:296
UInt256 previousLedgerHash_
Definition PeerImp.h:124
ClockType::time_point lastPingTime_
Definition PeerImp.h:131
json::Value json() override
Definition PeerImp.cpp:427
std::mutex recentLock_
Definition PeerImp.h:181
void onMessageEnd(std::uint16_t type, std::shared_ptr<::google::protobuf::Message > const &m)
Definition PeerImp.cpp:1088
std::uint32_t ID
Uniquely identifies a peer.
A public key.
Definition PublicKey.h:53
Slice slice() const noexcept
Definition PublicKey.h:115
A peer's signed, proposed position for use in RCLConsensus.
PublicKey const & publicKey() const
Public key of peer that sent the proposal.
ConsensusProposal< NodeID, UInt256, UInt256 > Proposal
UInt256 const & suppressionID() const
Unique id used by hash router to suppress duplicates.
bool checkSign() const
Verify the signing hash of the proposal.
bool getNodeFat(SHAMapNodeID const &wanted, std::vector< SHAMapNodeData > &data, bool fatLeaves, std::uint32_t depth) const
SHAMapHash getHash() const
void const * getDataPtr() const
Definition Serializer.h:294
int getLength() const
Definition Serializer.h:304
std::size_t size() const noexcept
Definition Serializer.h:147
void const * data() const noexcept
Definition Serializer.h:153
An immutable linear range of bytes.
Definition Slice.h:28
static Category attribute(Category cat, IsFromCluster isFromCluster)
Limits a category to what the sender can report.
static Category categorize(::google::protobuf::Message const &message, protocol::MessageType type, bool inbound)
Given a protocol message, determine which traffic category it belongs to.
static time_point now()
static void sendValidatorList(Peer &peer, std::uint64_t peerSequence, PublicKey const &publisherKey, std::size_t maxSequence, std::uint32_t rawVersion, std::string const &rawManifest, std::map< std::size_t, ValidatorBlobInfo > const &blobInfos, HashRouter &hashRouter, beast::Journal j)
static std::vector< ValidatorBlobInfo > parseBlobs(std::uint32_t version, json::Value const &body)
Pull the blob/signature/manifest information out of the appropriate Json body fields depending on the...
A consumption charge.
Definition Charge.h:13
An endpoint that consumes resources.
Definition Consumer.h:20
T duration_cast(T... args)
T emplace_back(T... args)
T empty(T... args)
T end(T... args)
T find(T... args)
T for_each(T... args)
T get(T... args)
T lock(T... args)
T make_shared(T... args)
T max(T... args)
T min(T... args)
constexpr Zero kZero
Definition Zero.h:30
unsigned int UInt
@ Object
object value (collection of name/value pairs).
Definition json_value.h:29
STL namespace.
TER valid(STTx const &tx, ReadView const &view, AccountID const &src, beast::Journal j)
auto measureDurationAndLog(Func &&func, std::string_view actionDescription, std::chrono::duration< Rep, Period > maxDelay, beast::Journal const &journal)
Definition PerfLog.h:174
static constexpr std::size_t kMaxTxQueueSize
static constexpr auto kIdled
Charge const kFeeMalformedData
Charge const kFeeHeavyBurdenPeer
Charge const kFeeUselessData
Charge const kFeeRequestNoReply
Charge const kFeeMalformedRequest
Schedule of fees charged for imposing load on the server.
Charge const kFeeInvalidSignature
Charge const kFeeInvalidData
Charge const kFeeModerateBurdenPeer
Charge const kFeeTrivialPeer
static constexpr auto kSoftMaxReplyNodes
The soft cap on the number of ledger entries in a single reply.
static constexpr auto kBandMediumMax
static constexpr auto kTargetSendQueue
How many messages we consider reasonable sustained on a send queue.
static constexpr auto kCostPerLookupHit
Cost of one cache-hit lookup.
static constexpr auto kCostBandLarge
static constexpr auto kFreeObjectsPerRequest
TMGetObjectByHash differential pricing.
static constexpr auto kHardMaxReplyNodes
The hard cap on the number of ledger entries in a single reply.
static constexpr std::uint32_t kConvergedLedgerLimit
How many ledgers off a server can be and we will still consider it converged.
static constexpr auto kDropSendQueue
How many messages on a send queue before we refuse queries.
static constexpr auto kMaxQueryDepth
The maximum number of levels to search.
static constexpr auto kSendqIntervals
How many timer intervals a sendq has to stay large before we disconnect.
static constexpr auto kBandSmallMax
Cutoffs that decide which size band a request falls into.
constexpr std::size_t kReadBufferBytes
Size of buffer used to read from the socket.
static constexpr auto kCostBandMedium
static constexpr auto kCostPerLookupMiss
Cost of one node-store miss, in units of kCostPerLookupHit.
static constexpr std::uint32_t kDivergedLedgerLimit
How many ledgers off a server has to be before we consider it diverged.
static constexpr auto kCostBandSmall
Size-band surcharges.
static constexpr auto kSendQueueLogFreq
How often to log send queue size.
Use hash_* containers for keys that do not need a cryptographically secure hashing algorithm.
Definition algorithm.h:5
bool set(T &target, std::string const &name, Section const &section)
Set a value from a configuration Section If the named value is not found or doesn't parse as a T,...
std::string base64Decode(std::string_view data)
constexpr FlagValue tfInnerBatchTxn
Definition TxFlags.h:44
bool isCurrent(ValidationParms const &p, NetClock::time_point now, NetClock::time_point signTime, NetClock::time_point seenTime)
Whether a validation is still current.
@ Valid
Signature and local checks are good / passed.
Definition apply.h:36
std::optional< AccountID > parseBase58(std::string const &s)
Parse AccountID from checked, base58 string.
std::uint32_t LedgerIndex
A ledger index.
Definition Protocol.h:382
Stopwatch & stopwatch()
Returns an instance of a wall clock.
Definition chrono.h:101
std::string strHex(FwdIt begin, FwdIt end)
Definition strHex.h:13
std::pair< Validity, std::string > checkValidity(HashRouter &router, STTx const &tx, Rules const &rules)
Checks transaction signature and local checks.
Definition apply.cpp:62
HttpResponseType makeResponse(bool crawlPublic, HttpRequestType const &req, beast::ip::Address publicIp, beast::ip::Address remoteIp, UInt256 const &sharedValue, std::optional< std::uint32_t > networkID, ProtocolVersion protocol, Application &app)
Make http response.
UInt256 proposalUniqueId(UInt256 const &proposeHash, UInt256 const &previousLedger, std::uint32_t proposeSeq, NetClock::time_point closeTime, Slice const &publicKey, Slice const &signature)
Calculate a unique identifier for a signed proposal.
@ UnsupportedVersion
List version is not supported.
@ Expired
List is expired, but has the largest non-pending sequence seen so far.
@ SameSequence
Same sequence as current list.
@ KnownSequence
Future sequence already seen.
@ Pending
List will be valid in the future.
@ Accepted
List is valid.
@ Invalid
Invalid format or signature.
@ Untrusted
List signed by untrusted publisher key.
@ Stale
Trusted publisher key, but seq is too old.
std::string toBase58(AccountID const &v)
Convert AccountID to base58 checked string.
Definition AccountID.cpp:95
static constexpr char kFeatureLedgerReplay[]
Definition Handshake.h:131
Number root(Number f, unsigned d)
static std::shared_ptr< PeerImp > getPeerWithTree(OverlayImpl &ov, UInt256 const &rootHash, PeerImp const *skip)
Definition PeerImp.cpp:3264
constexpr Dest safeCast(Src s) noexcept
Definition safe_cast.h:21
std::optional< SHAMapNodeID > getSHAMapNodeID(protocol::TMLedgerNode const &ledgerNode, SHAMapTreeNode const &treeNode)
Extracts or reconstructs the SHAMapNodeID from a ledger node proto message.
std::string to_string(BaseUInt< Bits, Tag > const &a)
Definition base_uint.h:657
std::optional< KeyType > publicKeyType(Slice const &slice)
Returns the type of public key.
SHAMapTreeNodePtr getTreeNode(std::string_view data)
Deserializes a SHAMapTreeNode from wire format data.
static constexpr char kFeatureTxrr[]
Definition Handshake.h:129
Slice makeSlice(std::array< T, N > const &a)
Definition Slice.h:228
std::pair< std::size_t, boost::system::error_code > invokeProtocolMessage(Buffers const &buffers, Handler &handler, std::size_t &hint)
Calls the handler for up to one protocol message in the passed buffers.
void forceValidity(HashRouter &router, UInt256 const &txid, Validity validity)
Sets the validity of a given transaction in the cache.
Definition apply.cpp:141
BaseUInt< 256 > UInt256
Definition base_uint.h:580
@ JtPack
Definition Job.h:29
@ JtLedgerReq
Definition Job.h:45
@ JtValidationUt
Definition Job.h:40
@ JtTransaction
Definition Job.h:48
@ JtProposalUt
Definition Job.h:46
@ JtReplayReq
Definition Job.h:44
@ JtRequestedTxn
Definition Job.h:50
@ JtTxnData
Definition Job.h:55
@ JtValidationT
Definition Job.h:57
@ JtMissingTxn
Definition Job.h:49
@ JtManifest
Definition Job.h:41
@ JtProposalT
Definition Job.h:60
@ JtPeer
Definition Job.h:66
void addRaw(LedgerHeader const &, Serializer &, bool includeHash=false)
boost::beast::http::request< boost::beast::http::dynamic_body > HttpRequestType
Definition Handoff.h:12
NodeID calcNodeID(PublicKey const &)
Calculate the 160-bit node ID from a node public key.
std::optional< UInt256 > makeSharedValue(StreamType &ssl, beast::Journal journal)
Computes a shared value based on the SSL connection state.
constexpr ProtocolVersion makeProtocol(std::uint16_t major, std::uint16_t minor)
HashRouterFlags
Definition HashRouter.h:20
std::optional< SHAMapNodeID > deserializeSHAMapNodeID(void const *data, std::size_t size)
Return an object representing a serialized SHAMap Node ID.
std::string getFingerprint(beast::ip::Endpoint const &address, std::optional< PublicKey > const &publicKey=std::nullopt, std::optional< std::string > const &id=std::nullopt)
Definition PublicKey.h:262
bool peerFeatureEnabled(Headers const &request, std::string const &feature, std::string value, bool config)
Check if a feature should be enabled for a peer.
Definition Handshake.h:182
@ Malformed
Protocol-level violation; no honest peer would produce this.
@ BadData
Peer reported has_error() (legitimate "cannot fulfill" signal).
std::pair< std::uint16_t, std::uint16_t > ProtocolVersion
Represents a particular version of the peer-to-peer protocol.
static bool stringIsUInt256Sized(std::string const &pBuffStr)
Definition PeerImp.cpp:200
static constexpr char kFeatureCompr[]
Definition Handshake.h:125
static std::shared_ptr< PeerImp > getPeerWithLedger(OverlayImpl &ov, UInt256 const &ledgerHash, LedgerIndex ledger, PeerImp const *skip)
Definition PeerImp.cpp:3288
std::string protocolMessageName(int type)
Returns the name of a protocol message given its type.
bool isPseudoTx(STObject const &tx)
Check whether a transaction is a pseudo-transaction.
Definition STTx.cpp:889
static constexpr char kFeatureVprr[]
Definition Handshake.h:127
Sha512HalfHasher::result_type sha512Half(Args const &... args)
Returns the SHA512-Half of a series of objects.
Definition digest.h:215
constexpr bool any(HashRouterFlags flags)
Definition HashRouter.h:74
T nth_element(T... args)
T push_back(T... args)
T ref(T... args)
T reserve(T... args)
T reset(T... args)
T size(T... args)
T str(T... args)
Information about the notional ledger backing the view.
Options controlling deserialization of a STValidation.
Describes a single consumer.
Definition Gossip.h:20
beast::ip::Endpoint address
Definition Gossip.h:24
Data format for exchanging consumption information across peers.
Definition Gossip.h:13
std::vector< Item > items
Definition Gossip.h:27
T tie(T... args)
T to_string(T... args)
T what(T... args)