xrpld
Loading...
Searching...
No Matches
OverlayImpl.cpp
1#include <xrpld/overlay/detail/OverlayImpl.h>
2
3#include <xrpld/app/misc/ValidatorList.h>
4#include <xrpld/app/misc/ValidatorSite.h>
5#include <xrpld/overlay/Cluster.h>
6#include <xrpld/overlay/detail/ConnectAttempt.h>
7#include <xrpld/overlay/detail/Handshake.h>
8#include <xrpld/overlay/detail/PeerImp.h>
9#include <xrpld/overlay/detail/ProtocolVersion.h>
10#include <xrpld/overlay/detail/TrafficCount.h>
11#include <xrpld/overlay/detail/Tuning.h>
12#include <xrpld/peerfinder/PeerfinderManager.h>
13#include <xrpld/rpc/ServerHandler.h>
14#include <xrpld/rpc/handlers/admin/status/GetCounts.h>
15#include <xrpld/rpc/json_body.h>
16
17#include <xrpl/basics/Log.h>
18#include <xrpl/basics/Resolver.h>
19#include <xrpl/basics/Slice.h>
20#include <xrpl/basics/base64.h>
21#include <xrpl/basics/base_uint.h>
22#include <xrpl/basics/chrono.h>
23#include <xrpl/basics/contract.h>
24#include <xrpl/basics/make_SSLContext.h>
25#include <xrpl/basics/random.h>
26#include <xrpl/basics/strHex.h>
27#include <xrpl/beast/core/LexicalCast.h>
28#include <xrpl/beast/insight/Collector.h>
29#include <xrpl/beast/net/IPAddress.h>
30#include <xrpl/beast/net/IPAddressConversion.h>
31#include <xrpl/beast/net/IPEndpoint.h>
32#include <xrpl/beast/rfc2616.h>
33#include <xrpl/beast/utility/PropertyStream.h>
34#include <xrpl/beast/utility/WrappedSink.h>
35#include <xrpl/beast/utility/instrumentation.h>
36#include <xrpl/config/BasicConfig.h>
37#include <xrpl/config/Constants.h>
38#include <xrpl/core/HashRouter.h>
39#include <xrpl/json/json_value.h>
40#include <xrpl/peerfinder/Config.h>
41#include <xrpl/peerfinder/Slot.h>
42#include <xrpl/peerfinder/make_Manager.h>
43#include <xrpl/protocol/BuildInfo.h>
44#include <xrpl/protocol/STTx.h>
45#include <xrpl/protocol/Serializer.h>
46#include <xrpl/protocol/SystemParameters.h>
47#include <xrpl/protocol/jss.h>
48#include <xrpl/resource/Fees.h>
49#include <xrpl/resource/ResourceManager.h>
50#include <xrpl/server/Handoff.h>
51#include <xrpl/server/Manifest.h>
52#include <xrpl/server/NetworkOPs.h>
53#include <xrpl/server/SimpleWriter.h>
54#include <xrpl/server/Wallet.h>
55#include <xrpl/server/Writer.h>
56
57#include <boost/algorithm/string/predicate.hpp>
58#include <boost/asio/bind_executor.hpp>
59#include <boost/asio/dispatch.hpp>
60#include <boost/asio/error.hpp>
61#include <boost/asio/executor_work_guard.hpp>
62#include <boost/asio/io_context.hpp>
63#include <boost/asio/ip/address.hpp>
64#include <boost/asio/post.hpp>
65#include <boost/asio/strand.hpp>
66#include <boost/beast/http/empty_body.hpp>
67#include <boost/beast/http/field.hpp>
68#include <boost/beast/http/status.hpp>
69#include <boost/lexical_cast.hpp>
70#include <boost/lexical_cast/bad_lexical_cast.hpp>
71#include <boost/lexical_cast/try_lexical_convert.hpp>
72
73#include <xrpl.pb.h>
74
75#include <algorithm>
76#include <chrono>
77#include <cstddef>
78#include <cstdint>
79#include <exception>
80#include <format>
81#include <functional>
82#include <iomanip>
83#include <memory>
84#include <mutex>
85#include <optional>
86#include <set>
87#include <sstream>
88#include <stdexcept>
89#include <string>
90#include <string_view>
91#include <tuple>
92#include <unordered_map>
93#include <utility>
94#include <vector>
95
96namespace xrpl {
97
98namespace crawl_options {
99static constexpr auto kDisabled = 0;
100static constexpr auto kOverlay = (1 << 0);
101static constexpr auto kServerInfo = (1 << 1);
102static constexpr auto kServerCounts = (1 << 2);
103static constexpr auto kUnl = (1 << 3);
104} // namespace crawl_options
105
106//------------------------------------------------------------------------------
107
109{
110}
111
113{
114 overlay_.remove(*this);
115}
116
117//------------------------------------------------------------------------------
118
122
123void
125{
126 // This method is only ever called from the same strand that calls
127 // Timer::on_timer, ensuring they never execute concurrently.
128 stopping = true;
129 timer.cancel();
130}
131
132void
134{
135 timer.expires_after(std::chrono::seconds(1));
136 timer.async_wait(
137 boost::asio::bind_executor(
138 overlay_.strand_,
139 [self = shared_from_this()](ErrorCode const& ec) { self->onTimer(ec); }));
140}
141
142void
144{
145 if (ec || stopping)
146 {
147 if (ec && ec != boost::asio::error::operation_aborted)
148 {
149 JLOG(overlay_.journal_.error()) << "on_timer: " << ec.message();
150 }
151 return;
152 }
153
154 overlay_.peerFinder_->oncePerSecond();
155 overlay_.sendEndpoints();
156 overlay_.autoConnect();
157 if (overlay_.app_.config().txReduceRelayEnable)
158 overlay_.sendTxQueue();
159
160 if ((++overlay_.timerCount_ % tuning::kCheckIdlePeers) == 0)
161 overlay_.deleteIdlePeers();
162
163 asyncWait();
164}
165
166//------------------------------------------------------------------------------
167
169 Application& app,
170 Setup setup,
171 ServerHandler& serverHandler,
173 Resolver& resolver,
174 boost::asio::io_context& ioContext,
175 BasicConfig const& config,
176 beast::insight::Collector::Ptr const& collector)
177 : app_(app)
178 , ioContext_(ioContext)
179 , work_(std::in_place, boost::asio::make_work_guard(ioContext_))
180 , strand_(boost::asio::make_strand(ioContext_))
181 , setup_(std::move(setup))
182 , journal_(app_.getJournal("Overlay"))
183 , serverHandler_(serverHandler)
185 , store_(app_.getJournal("PeerFinder"))
186 , peerFinder_(
187 peer_finder::makeManager(
188 ioContext,
189 stopwatch(),
190 app_.getJournal("PeerFinder"),
191 store_,
192 collector))
193 , resolver_(resolver)
194 , nextId_(1)
195 , slots_(app, *this, app.config())
196 , stats_(
197 [this] { collectMetrics(); },
198 collector,
199 [counts = traffic_.getCounts(), collector]() {
201
202 for (auto const& pair : counts)
203 ret.emplace(pair.first, TrafficGauges(pair.second.name, collector));
204
205 return ret;
206 }())
207{
208 store_.open(config);
209 beast::PropertyStream::Source::add(peerFinder_.get());
210}
211
214 std::unique_ptr<StreamType>&& streamPtr,
215 HttpRequestType&& request,
216 EndpointType remoteEndpoint)
217{
218 auto const id = nextId_++;
219 auto peerJournal = app_.getJournal("Peer");
220 beast::WrappedSink sink(peerJournal.sink(), makePrefix(id));
221 beast::Journal const journal(sink);
222
223 Handoff handoff;
224 if (processRequest(request, handoff))
225 return handoff;
226 if (!isPeerUpgrade(request))
227 return handoff;
228
229 handoff.moved = true;
230
231 JLOG(journal.debug()) << "Peer connection upgrade from " << remoteEndpoint;
232
233 ErrorCode ec;
234 auto const localEndpoint(streamPtr->next_layer().socket().local_endpoint(ec));
235 if (ec)
236 {
237 JLOG(journal.debug()) << remoteEndpoint << " failed: " << ec.message();
238 return handoff;
239 }
240
241 auto consumer =
242 resourceManager_.newInboundEndpoint(beast::IPAddressConversion::fromAsio(remoteEndpoint));
243 if (consumer.disconnect(journal))
244 return handoff;
245
246 auto const [slot, result] = peerFinder_->newInboundSlot(
249
250 if (slot == nullptr)
251 {
252 // connection refused either IP limit exceeded or self-connect
253 handoff.moved = false;
254 JLOG(journal.debug()) << "Peer " << remoteEndpoint << " refused, " << to_string(result);
255 return handoff;
256 }
257
258 // Validate HTTP request
259
260 {
261 auto const types = beast::rfc2616::splitCommas(request["Connect-As"]);
262 if (std::ranges::find_if(types, [](std::string const& s) {
263 return boost::iequals(s, "peer");
264 }) == types.end())
265 {
266 handoff.moved = false;
267 handoff.response = makeRedirectResponse(slot, request, remoteEndpoint.address());
268 handoff.keepAlive = beast::rfc2616::isKeepAlive(request);
269 return handoff;
270 }
271 }
272
273 auto const negotiatedVersion = negotiateProtocolVersion(request["Upgrade"]);
274 if (!negotiatedVersion)
275 {
276 peerFinder_->onClosed(slot);
277 handoff.moved = false;
278 handoff.response = makeErrorResponse(
279 slot, request, remoteEndpoint.address(), "Unable to agree on a protocol version");
280 handoff.keepAlive = false;
281 return handoff;
282 }
283
284 auto const sharedValue = makeSharedValue(*streamPtr, journal);
285 if (!sharedValue)
286 {
287 peerFinder_->onClosed(slot);
288 handoff.moved = false;
289 handoff.response =
290 makeErrorResponse(slot, request, remoteEndpoint.address(), "Incorrect security cookie");
291 handoff.keepAlive = false;
292 return handoff;
293 }
294
295 try
296 {
297 auto publicKey = verifyHandshake(
298 request,
299 *sharedValue,
300 setup_.networkID,
301 setup_.publicIp,
302 remoteEndpoint.address(),
303 app_);
304
305 consumer.setPublicKey(publicKey);
306
307 {
308 // The node gets a reserved slot if it is in our cluster
309 // or if it has a reservation.
310 bool const reserved = app_.getCluster().isMember(publicKey) ||
311 app_.getPeerReservations().contains(publicKey);
312 auto const result = peerFinder_->activate(slot, publicKey, reserved);
313 if (result != peer_finder::Result::Success)
314 {
315 peerFinder_->onClosed(slot);
316 JLOG(journal.debug())
317 << "Peer " << remoteEndpoint << " redirected, " << to_string(result);
318 handoff.moved = false;
319 handoff.response = makeRedirectResponse(slot, request, remoteEndpoint.address());
320 handoff.keepAlive = false;
321 return handoff;
322 }
323 }
324
325 auto const peer = std::make_shared<PeerImp>(
326 app_,
327 id,
328 slot,
329 std::move(request),
330 publicKey,
331 *negotiatedVersion,
332 consumer,
333 std::move(streamPtr),
334 *this);
335 {
336 // As we are not on the strand, run() must be called
337 // while holding the lock, otherwise new I/O can be
338 // queued after a call to stop().
339 std::scoped_lock const lock(mutex_);
340 {
341 auto const result = peers_.emplace(peer->slot(), peer);
342 XRPL_ASSERT(result.second, "xrpl::OverlayImpl::onHandoff : peer is inserted");
343 (void)result.second;
344 }
345 list_.emplace(peer.get(), peer);
346
347 peer->run();
348 }
349 handoff.moved = true;
350 return handoff;
351 }
352 catch (std::exception const& e)
353 {
354 JLOG(journal.debug()) << "Peer " << remoteEndpoint << " fails handshake (" << e.what()
355 << ")";
356
357 peerFinder_->onClosed(slot);
358 handoff.moved = false;
359 handoff.response = makeErrorResponse(slot, request, remoteEndpoint.address(), e.what());
360 handoff.keepAlive = false;
361 return handoff;
362 }
363}
364
365//------------------------------------------------------------------------------
366
367bool
369{
370 if (!isUpgrade(request))
371 return false;
372 auto const versions = parseProtocolVersions(request["Upgrade"]);
373 return !versions.empty();
374}
375
378{
380 ss << "[" << std::setfill('0') << std::setw(3) << id << "] ";
381 return ss.str();
382}
383
387 HttpRequestType const& request,
388 AddressType remoteAddress)
389{
390 boost::beast::http::response<JsonBody> msg;
391 msg.version(request.version());
392 msg.result(boost::beast::http::status::service_unavailable);
393 msg.insert("Server", build_info::getFullVersionString());
394 {
396 ostr << remoteAddress;
397 msg.insert("Remote-Address", ostr.str());
398 }
399 msg.insert("Content-Type", "application/json");
400 msg.insert(boost::beast::http::field::connection, "close");
401 msg.body() = json::ValueType::Object;
402 {
403 json::Value& ips = (msg.body()["peer-ips"] = json::ValueType::Array);
404 for (auto const& _ : peerFinder_->redirect(slot))
405 ips.append(_.address.toString());
406 }
407 msg.prepare_payload();
409}
410
414 HttpRequestType const& request,
415 AddressType remoteAddress,
416 std::string const& text)
417{
418 boost::beast::http::response<boost::beast::http::empty_body> msg;
419 msg.version(request.version());
420 msg.result(boost::beast::http::status::bad_request);
421 msg.reason("Bad Request (" + text + ")");
422 msg.insert("Server", build_info::getFullVersionString());
423 msg.insert("Remote-Address", remoteAddress.to_string());
424 msg.insert(boost::beast::http::field::connection, "close");
425 msg.prepare_payload();
427}
428
429//------------------------------------------------------------------------------
430
431void
433{
434 XRPL_ASSERT(work_, "xrpl::OverlayImpl::connect : work is set");
435
436 auto usage = resourceManager().newOutboundEndpoint(remoteEndpoint);
437 if (usage.disconnect(journal_))
438 {
439 JLOG(journal_.info()) << "Over resource limit: " << remoteEndpoint;
440 return;
441 }
442
443 auto const [slot, result] = peerFinder().newOutboundSlot(remoteEndpoint);
444 if (slot == nullptr)
445 {
446 JLOG(journal_.debug()) << "Connect: No slot for " << remoteEndpoint << ": "
447 << to_string(result);
448 return;
449 }
450
452 app_,
455 usage,
456 setup_.context,
457 nextId_++,
458 slot,
459 app_.getJournal("Peer"),
460 *this);
461
462 std::scoped_lock const lock(mutex_);
463 list_.emplace(p.get(), p);
464 p->run();
465}
466
467//------------------------------------------------------------------------------
468
469// Adds a peer that is already handshaked and active
470void
472{
473 beast::WrappedSink sink{journal_.sink(), peer->prefix()};
474 beast::Journal const journal{sink};
475
476 std::scoped_lock const lock(mutex_);
477
478 {
479 auto const result = peers_.emplace(peer->slot(), peer);
480 XRPL_ASSERT(result.second, "xrpl::OverlayImpl::addActive : peer is inserted");
481 (void)result.second;
482 }
483
484 {
485 auto const result = ids_.emplace(
487 XRPL_ASSERT(result.second, "xrpl::OverlayImpl::addActive : peer ID is inserted");
488 (void)result.second;
489 }
490
491 list_.emplace(peer.get(), peer);
492
493 JLOG(journal.debug()) << "activated";
494
495 // As we are not on the strand, run() must be called
496 // while holding the lock, otherwise new I/O can be
497 // queued after a call to stop().
498 peer->run();
499}
500
501void
503{
504 std::scoped_lock const lock(mutex_);
505 auto const iter = peers_.find(slot);
506 XRPL_ASSERT(iter != peers_.end(), "xrpl::OverlayImpl::remove : valid input");
507 peers_.erase(iter);
508}
509
510void
512{
514 app_.config(),
515 serverHandler_.setup().overlay.port(),
516 app_.getValidationPublicKey().has_value(),
517 setup_.ipLimit,
518 setup_.verifyEndpoints);
519
520 peerFinder_->setConfig(config);
521 peerFinder_->start();
522
523 // Populate our boot cache: if there are no entries in [ips] then we use
524 // the entries in [ips_fixed].
525 auto bootstrapIps = app_.config().ips.empty() ? app_.config().ipsFixed : app_.config().ips;
526
527 // If nothing is specified, default to several well-known high-capacity
528 // servers to serve as bootstrap:
529 if (bootstrapIps.empty())
530 {
531 // Pool of servers operated by Ripple Labs Inc. - https://ripple.com
532 bootstrapIps.emplace_back("r.ripple.com 51235");
533
534 // Pool of servers operated by ISRDC - https://isrdc.in
535 bootstrapIps.emplace_back("sahyadri.isrdc.in 51235");
536
537 // Pool of servers operated by @Xrpkuwait - https://xrpkuwait.com
538 bootstrapIps.emplace_back("hubs.xrpkuwait.com 51235");
539
540 // Pool of servers operated by XRPL Commons - https://xrpl-commons.org
541 bootstrapIps.emplace_back("hub.xrpl-commons.org 51235");
542 }
543
544 resolver_.resolve(
545 bootstrapIps,
546 [this](std::string const& name, std::vector<beast::ip::Endpoint> const& addresses) {
548 ips.reserve(addresses.size());
549 for (auto const& addr : addresses)
550 {
551 if (addr.port() == 0)
552 {
553 ips.push_back(to_string(addr.atPort(kDefaultPeerPort)));
554 }
555 else
556 {
557 ips.push_back(to_string(addr));
558 }
559 }
560
561 std::string const base("config: ");
562 if (!ips.empty())
563 peerFinder_->addFallbackStrings(base + name, ips);
564 });
565
566 // Add the ips_fixed from the xrpld.cfg file
567 if (!app_.config().standalone() && !app_.config().ipsFixed.empty())
568 {
569 resolver_.resolve(
570 app_.config().ipsFixed,
571 [this](std::string const& name, std::vector<beast::ip::Endpoint> const& addresses) {
572 std::vector<beast::ip::Endpoint> ips;
573 ips.reserve(addresses.size());
574
575 for (auto& addr : addresses)
576 {
577 if (addr.port() == 0)
578 {
579 ips.emplace_back(addr.address(), kDefaultPeerPort);
580 }
581 else
582 {
583 ips.emplace_back(addr);
584 }
585 }
586
587 if (!ips.empty())
588 peerFinder_->addFixedPeer(name, ips);
589 });
590 }
591 auto const timer = std::make_shared<Timer>(*this);
592 std::scoped_lock const lock(mutex_);
593 list_.emplace(timer.get(), timer);
594 timer_ = timer;
595 timer->asyncWait();
596}
597
598void
600{
601 boost::asio::dispatch(strand_, [this] { stopChildren(); });
602 {
603 std::unique_lock<decltype(mutex_)> lock(mutex_);
604 cond_.wait(lock, [this] { return list_.empty(); });
605 }
606 peerFinder_->stop();
607}
608
609//------------------------------------------------------------------------------
610//
611// PropertyStream
612//
613//------------------------------------------------------------------------------
614
615void
617{
618 beast::PropertyStream::Set set("traffic", stream);
619 auto const stats = traffic_.getCounts();
620 for (auto const& pair : stats)
621 {
623 item["category"] = pair.second.name;
624 item["bytes_in"] = std::to_string(pair.second.bytesIn.load());
625 item["messages_in"] = std::to_string(pair.second.messagesIn.load());
626 item["bytes_out"] = std::to_string(pair.second.bytesOut.load());
627 item["messages_out"] = std::to_string(pair.second.messagesOut.load());
628 }
629}
630
631//------------------------------------------------------------------------------
638void
640{
641 beast::WrappedSink sink{journal_.sink(), peer->prefix()};
642 beast::Journal const journal{sink};
643
644 // Now track this peer
645 {
646 std::scoped_lock const lock(mutex_);
647 auto const result(ids_.emplace(
649 XRPL_ASSERT(result.second, "xrpl::OverlayImpl::activate : peer ID is inserted");
650 (void)result.second;
651 }
652
653 JLOG(journal.debug()) << "activated";
654
655 // We just accepted this peer so we have non-zero active peers
656 XRPL_ASSERT(size(), "xrpl::OverlayImpl::activate : nonzero peers");
657}
658
659void
661{
662 std::scoped_lock const lock(mutex_);
663 ids_.erase(id);
664}
665
666void
669 std::shared_ptr<PeerImp> const& from)
670{
671 auto const& journal = from->pJournal();
672
673 // Process every trusted manifest, but stop processing untrusted ones once
674 // the configured untrusted count has been handled, so the work stays
675 // bounded. Trusted manifests are always processed: dropping one would delay
676 // a validator key rotation reaching this node.
677 auto const maxUntrusted = untrustedManifestCount(app_.config().maxUntrustedCount);
678 auto const total = static_cast<std::size_t>(m->list_size());
679 std::size_t untrusted = 0;
680 bool skippedUntrusted = false;
681
682 protocol::TMManifests relay;
683
684 for (std::size_t i = 0; i < total; ++i)
685 {
686 auto& s = m->list().Get(i).stobject();
687
688 if (auto mo = deserializeManifest(s))
689 {
690 auto const serialized = mo->serialized;
691 // Resolve trust before applyManifest takes the manifest-cache
692 // lock: listed() takes the validator-list lock, so ordering it
693 // first avoids holding the two locks in opposite orders.
694 bool const isTrusted = app_.getValidators().listed(mo->masterKey);
695
696 // Bound untrusted work: process at most maxUntrusted untrusted
697 // manifests, but never skip a trusted one. Trusted manifests are
698 // not counted against the cap.
699 if (!isTrusted)
700 {
701 if (untrusted >= maxUntrusted)
702 {
703 skippedUntrusted = true;
704 continue;
705 }
706 ++untrusted;
707 }
708
709 auto const result = app_.getValidatorManifests().applyManifest(
710 std::move(*mo),
713
714 if (result == ManifestDisposition::Accepted)
715 {
716 // N.B.: this is important; the applyManifest call above moves
717 // the loaded Manifest out of the optional so we need to
718 // reload it here.
719 mo = deserializeManifest(serialized);
720 XRPL_ASSERT(
721 mo,
722 "xrpl::OverlayImpl::onManifests : manifest "
723 "deserialization succeeded");
724 // NOLINTBEGIN(bugprone-unchecked-optional-access) assert above
725 app_.getOPs().pubManifest(*mo);
726 // NOLINTEND(bugprone-unchecked-optional-access)
727
728 relay.add_list()->set_stobject(s);
729
730 // Persist to the wallet DB only for trusted keys, so untrusted
731 // gossip never survives a restart.
732 if (isTrusted)
733 {
734 auto db = app_.getWalletDB().checkoutDb();
735 addValidatorManifest(*db, serialized);
736 }
737 }
738 }
739 else
740 {
741 JLOG(journal.debug()) << "Malformed manifest #" << i + 1 << ": " << strHex(s);
742 continue;
743 }
744 }
745
746 if (skippedUntrusted)
747 {
748 // The sender exceeded the untrusted per-message cap. Charge it (once,
749 // here) so a flood of untrusted manifests is penalized.
750 from->charge(resource::kFeeMalformedRequest, "too many untrusted manifests");
751
752 JLOG(journal.warn()) << "Manifests: message had " << total
753 << " entries; processed all trusted plus the first " << maxUntrusted
754 << " untrusted";
755 }
756
757 if (!relay.list().empty())
758 {
759 forEach([m2 = std::make_shared<Message>(relay, protocol::mtMANIFESTS)](
760 std::shared_ptr<PeerImp> const& p) { p->send(m2); });
761 }
762}
763
764void
766{
767 traffic_.addCount(cat, true, size);
768}
769
770void
772{
773 traffic_.addCount(cat, false, size);
774}
775
782{
783 std::scoped_lock const lock(mutex_);
784 return ids_.size();
785}
786
787int
789{
790 return peerFinder_->config().maxPeers;
791}
792
795{
796 using namespace std::chrono;
797 json::Value jv;
798 auto& av = jv[jss::active] = json::Value(json::ValueType::Array);
799
800 forEach([&](std::shared_ptr<PeerImp> const& sp) {
801 auto& pv = av.append(json::Value(json::ValueType::Object));
802 pv[jss::public_key] = base64Encode(sp->getNodePublic().data(), sp->getNodePublic().size());
803 pv[jss::type] = sp->slot()->inbound() ? jss::in : jss::out;
804 pv[jss::uptime] = static_cast<std::uint32_t>(duration_cast<seconds>(sp->uptime()).count());
805 if (sp->crawl())
806 {
807 pv[jss::ip] = sp->getRemoteAddress().address().to_string();
808 if (sp->slot()->inbound())
809 {
810 if (auto port = sp->slot()->listeningPort())
811 pv[jss::port] = *port;
812 }
813 else
814 {
815 pv[jss::port] = sp->getRemoteAddress().port();
816 }
817 }
818
819 {
820 auto version{sp->getVersion()};
821 if (!version.empty())
822 {
823 // Could move here if json::value supported moving from strings
824 pv[jss::version] = std::string{version};
825 }
826 }
827
828 std::uint32_t minSeq = 0, maxSeq = 0;
829 sp->ledgerRange(minSeq, maxSeq);
830 if (minSeq != 0 || maxSeq != 0)
831 pv[jss::complete_ledgers] = std::to_string(minSeq) + "-" + std::to_string(maxSeq);
832 });
833
834 return jv;
835}
836
839{
840 bool const humanReadable = false;
841 bool const admin = false;
842 bool const counters = false;
843
844 json::Value serverInfo = app_.getOPs().getServerInfo(humanReadable, admin, counters);
845
846 // Filter out some information
847 serverInfo.removeMember(jss::hostid);
848 serverInfo.removeMember(jss::load_factor_fee_escalation);
849 serverInfo.removeMember(jss::load_factor_fee_queue);
850 serverInfo.removeMember(jss::validation_quorum);
851
852 if (serverInfo.isMember(jss::validated_ledger))
853 {
854 json::Value& validatedLedger = serverInfo[jss::validated_ledger];
855
856 validatedLedger.removeMember(jss::base_fee);
857 validatedLedger.removeMember(jss::reserve_base_xrp);
858 validatedLedger.removeMember(jss::reserve_inc_xrp);
859 }
860
861 return serverInfo;
862}
863
869
872{
873 json::Value validators = app_.getValidators().getJson();
874
875 if (validators.isMember(jss::publisher_lists))
876 {
877 json::Value& publisherLists = validators[jss::publisher_lists];
878
879 for (auto& publisher : publisherLists)
880 {
881 publisher.removeMember(jss::list);
882 }
883 }
884
885 validators.removeMember(jss::signing_keys);
886 validators.removeMember(jss::trusted_validator_keys);
887 validators.removeMember(jss::validation_quorum);
888
889 json::Value validatorSites = app_.getValidatorSites().getJson();
890
891 if (validatorSites.isMember(jss::validator_sites))
892 {
893 validators[jss::validator_sites] = std::move(validatorSites[jss::validator_sites]);
894 }
895
896 return validators;
897}
898
899// Returns information on verified peers.
902{
904 for (auto const& peer : getActivePeers())
905 {
906 json.append(peer->json());
907 }
908 return json;
909}
910
911bool
913{
914 if (req.target() != "/crawl" || setup_.crawlOptions == crawl_options::kDisabled)
915 return false;
916
917 boost::beast::http::response<JsonBody> msg;
918 msg.version(req.version());
919 msg.result(boost::beast::http::status::ok);
920 msg.insert("Server", build_info::getFullVersionString());
921 msg.insert("Content-Type", "application/json");
922 msg.insert("Connection", "close");
923 msg.body()["version"] = json::Value(2u);
924
925 if ((setup_.crawlOptions & crawl_options::kOverlay) != 0u)
926 {
927 msg.body()["overlay"] = getOverlayInfo();
928 }
929 if ((setup_.crawlOptions & crawl_options::kServerInfo) != 0u)
930 {
931 msg.body()["server"] = getServerInfo();
932 }
933 if ((setup_.crawlOptions & crawl_options::kServerCounts) != 0u)
934 {
935 msg.body()["counts"] = getServerCounts();
936 }
937 if ((setup_.crawlOptions & crawl_options::kUnl) != 0u)
938 {
939 msg.body()["unl"] = getUnlInfo();
940 }
941
942 msg.prepare_payload();
944 return true;
945}
946
947bool
949{
950 // If the target is in the form "/vl/<validator_list_public_key>",
951 // return the most recent validator list for that key.
952 constexpr std::string_view kPrefix("/vl/");
953
954 if (!req.target().starts_with(kPrefix) || !setup_.vlEnabled)
955 return false;
956
957 std::uint32_t version = 1;
958
959 boost::beast::http::response<JsonBody> msg;
960 msg.version(req.version());
961 msg.insert("Server", build_info::getFullVersionString());
962 msg.insert("Content-Type", "application/json");
963 msg.insert("Connection", "close");
964
965 auto fail = [&msg, &handoff](auto status) {
966 msg.result(status);
967 msg.insert("Content-Length", "0");
968
969 msg.body() = json::ValueType::Null;
970
971 msg.prepare_payload();
973 return true;
974 };
975
976 std::string_view key = req.target().substr(kPrefix.size());
977
978 if (auto slash = key.find('/'); slash != std::string_view::npos)
979 {
980 auto verString = key.substr(0, slash);
981 if (!boost::conversion::try_lexical_convert(verString, version))
982 return fail(boost::beast::http::status::bad_request);
983 key = key.substr(slash + 1);
984 }
985
986 if (key.empty())
987 return fail(boost::beast::http::status::bad_request);
988
989 // find the list
990 auto vl = app_.getValidators().getAvailable(key, version);
991
992 if (!vl)
993 {
994 // 404 not found
995 return fail(boost::beast::http::status::not_found);
996 }
997 if (!*vl)
998 {
999 return fail(boost::beast::http::status::bad_request);
1000 }
1001
1002 msg.result(boost::beast::http::status::ok);
1003
1004 msg.body() = *vl;
1005
1006 msg.prepare_payload();
1008 return true;
1009}
1010
1011bool
1013{
1014 if (req.target() != "/health")
1015 return false;
1016 boost::beast::http::response<JsonBody> msg;
1017 msg.version(req.version());
1018 msg.insert("Server", build_info::getFullVersionString());
1019 msg.insert("Content-Type", "application/json");
1020 msg.insert("Connection", "close");
1021
1022 auto info = getServerInfo();
1023
1024 int lastValidatedLedgerAge = -1;
1025 if (info.isMember(jss::validated_ledger))
1026 lastValidatedLedgerAge = info[jss::validated_ledger][jss::age].asInt();
1027 bool amendmentBlocked = false;
1028 if (info.isMember(jss::amendment_blocked))
1029 amendmentBlocked = true;
1030 int const numberPeers = info[jss::peers].asInt();
1031 std::string const serverState = info[jss::server_state].asString();
1032 auto loadFactor = info[jss::load_factor_server].asDouble() / info[jss::load_base].asDouble();
1033
1034 enum class HealthState { Healthy, Warning, Critical };
1035 auto health = HealthState::Healthy;
1036 auto setHealth = [&health](HealthState state) { health = std::max(health, state); };
1037
1038 msg.body()[jss::info] = json::ValueType::Object;
1039 if (lastValidatedLedgerAge >= 7 || lastValidatedLedgerAge < 0)
1040 {
1041 msg.body()[jss::info][jss::validated_ledger] = lastValidatedLedgerAge;
1042 if (lastValidatedLedgerAge < 20)
1043 {
1044 setHealth(HealthState::Warning);
1045 }
1046 else
1047 {
1048 setHealth(HealthState::Critical);
1049 }
1050 }
1051
1052 if (amendmentBlocked)
1053 {
1054 msg.body()[jss::info][jss::amendment_blocked] = true;
1055 setHealth(HealthState::Critical);
1056 }
1057
1058 if (numberPeers <= 7)
1059 {
1060 msg.body()[jss::info][jss::peers] = numberPeers;
1061 if (numberPeers != 0)
1062 {
1063 setHealth(HealthState::Warning);
1064 }
1065 else
1066 {
1067 setHealth(HealthState::Critical);
1068 }
1069 }
1070
1071 if (!(serverState == "full" || serverState == "validating" || serverState == "proposing"))
1072 {
1073 msg.body()[jss::info][jss::server_state] = serverState;
1074 if (serverState == "syncing" || serverState == "tracking" || serverState == "connected")
1075 {
1076 setHealth(HealthState::Warning);
1077 }
1078 else
1079 {
1080 setHealth(HealthState::Critical);
1081 }
1082 }
1083
1084 if (loadFactor > 100)
1085 {
1086 msg.body()[jss::info][jss::load_factor] = loadFactor;
1087 if (loadFactor < 1000)
1088 {
1089 setHealth(HealthState::Warning);
1090 }
1091 else
1092 {
1093 setHealth(HealthState::Critical);
1094 }
1095 }
1096
1097 switch (health)
1098 {
1099 case HealthState::Healthy:
1100 msg.result(boost::beast::http::status::ok);
1101 break;
1102 case HealthState::Warning:
1103 msg.result(boost::beast::http::status::service_unavailable);
1104 break;
1105 case HealthState::Critical:
1106 msg.result(boost::beast::http::status::internal_server_error);
1107 break;
1108 }
1109
1110 msg.prepare_payload();
1112 return true;
1113}
1114
1115bool
1117{
1118 // Take advantage of || short-circuiting
1119 return processCrawl(req, handoff) || processValidatorList(req, handoff) ||
1120 processHealth(req, handoff);
1121}
1122
1125{
1127 ret.reserve(size());
1128
1129 forEach([&ret](std::shared_ptr<PeerImp> const& sp) { ret.emplace_back(sp); });
1130
1131 return ret;
1132}
1133
1136 std::set<Peer::ID> const& toSkip,
1137 std::size_t& active,
1138 std::size_t& disabled,
1139 std::size_t& enabledInSkip) const
1140{
1143
1144 active = ids_.size();
1145 disabled = enabledInSkip = 0;
1146 ret.reserve(ids_.size());
1147
1148 // NOTE The purpose of p is to delay the destruction of PeerImp
1150 for (auto& [id, w] : ids_)
1151 {
1152 if (p = w.lock(); p != nullptr)
1153 {
1154 bool const reduceRelayEnabled = p->txReduceRelayEnabled();
1155 // tx reduced relay feature disabled
1156 if (!reduceRelayEnabled)
1157 ++disabled;
1158
1159 if (!toSkip.contains(id))
1160 {
1161 ret.emplace_back(std::move(p));
1162 }
1163 else if (reduceRelayEnabled)
1164 {
1165 ++enabledInSkip;
1166 }
1167 }
1168 }
1169
1170 return ret;
1171}
1172
1173void
1175{
1176 forEach([index](std::shared_ptr<PeerImp> const& sp) { sp->checkTracking(index); });
1177}
1178
1181{
1183 auto const iter = ids_.find(id);
1184 if (iter != ids_.end())
1185 return iter->second.lock();
1186 return {};
1187}
1188
1189// A public key hash map was not used due to the peer connect/disconnect
1190// update overhead outweighing the performance of a small set linear search.
1193{
1195 // NOTE The purpose of peer is to delay the destruction of PeerImp
1197 for (auto const& e : ids_)
1198 {
1199 if (peer = e.second.lock(); peer != nullptr)
1200 {
1201 if (peer->getNodePublic() == pubKey)
1202 return peer;
1203 }
1204 }
1205 return {};
1206}
1207
1208void
1209OverlayImpl::broadcast(protocol::TMProposeSet const& m)
1210{
1211 auto const sm = std::make_shared<Message>(m, protocol::mtPROPOSE_LEDGER);
1212 forEach([&](std::shared_ptr<PeerImp> const& p) { p->send(sm); });
1213}
1214
1216OverlayImpl::relay(protocol::TMProposeSet const& m, UInt256 const& uid, PublicKey const& validator)
1217{
1218 if (auto const toSkip = app_.getHashRouter().shouldRelay(uid))
1219 {
1220 auto const sm = std::make_shared<Message>(m, protocol::mtPROPOSE_LEDGER, validator);
1221 forEach([&](std::shared_ptr<PeerImp> const& p) {
1222 if (!toSkip->contains(p->id()))
1223 p->send(sm);
1224 });
1225 return *toSkip;
1226 }
1227 return {};
1228}
1229
1230void
1231OverlayImpl::broadcast(protocol::TMValidation const& m)
1232{
1233 auto const sm = std::make_shared<Message>(m, protocol::mtVALIDATION);
1234 forEach([sm](std::shared_ptr<PeerImp> const& p) { p->send(sm); });
1235}
1236
1238OverlayImpl::relay(protocol::TMValidation const& m, UInt256 const& uid, PublicKey const& validator)
1239{
1240 if (auto const toSkip = app_.getHashRouter().shouldRelay(uid))
1241 {
1242 auto const sm = std::make_shared<Message>(m, protocol::mtVALIDATION, validator);
1243 forEach([&](std::shared_ptr<PeerImp> const& p) {
1244 if (!toSkip->contains(p->id()))
1245 p->send(sm);
1246 });
1247 return *toSkip;
1248 }
1249 return {};
1250}
1251
1254{
1256
1257 if (auto seq = app_.getValidatorManifests().sequence(); seq != manifestListSeq_)
1258 {
1259 // Phase 1: snapshot the cache under its own lock. Do not call
1260 // Validators::listed() here — that takes the validator-list lock, and
1261 // forEachManifest holds the manifest-cache lock, so consulting trust
1262 // inside the callback would invert the lock order used elsewhere
1263 // (see onManifests) and risk deadlock. Capture the manifest hash now,
1264 // while we have the Manifest object, for the suppression key.
1265 struct CachedManifest
1266 {
1267 PublicKey masterKey;
1268 std::string serialized;
1269 UInt256 hash;
1270 };
1272 app_.getValidatorManifests().forEachManifest(
1273 [&cached](std::size_t s) { cached.reserve(s); },
1274 [&cached](Manifest const& manifest) {
1275 cached.push_back(
1276 {.masterKey = manifest.masterKey,
1277 .serialized = manifest.serialized,
1278 .hash = manifest.hash()});
1279 });
1280
1281 // Phase 2: no cache lock held, so trust checks are safe. Include every
1282 // trusted manifest, then fill any remaining headroom up to the
1283 // configured untrusted count with gossip. Trusted manifests are never
1284 // dropped; the trusted count only sizes the accepted message.
1287 for (auto const& e : cached)
1288 {
1289 if (app_.getValidators().listed(e.masterKey))
1290 {
1291 selected.push_back(&e);
1292 }
1293 else
1294 {
1295 untrusted.push_back(&e);
1296 }
1297 }
1298
1299 // Cap untrusted only; trusted manifests are all included above.
1300 auto const take =
1301 std::min(untrustedManifestCount(app_.config().maxUntrustedCount), untrusted.size());
1302 selected.insert(selected.end(), untrusted.begin(), untrusted.begin() + take);
1303
1304 // Shuffle the order. Cryptographic randomness is not needed here.
1305 std::shuffle(selected.begin(), selected.end(), defaultPrng());
1306
1307 protocol::TMManifests tm;
1308 auto& hr = app_.getHashRouter();
1309 tm.mutable_list()->Reserve(static_cast<int>(selected.size()));
1310 for (auto const* e : selected)
1311 {
1312 tm.add_list()->set_stobject(e->serialized.data(), e->serialized.size());
1313 hr.addSuppression(e->hash);
1314 }
1315
1316 manifestMessage_.reset();
1317
1318 if (tm.list_size() != 0)
1319 manifestMessage_ = std::make_shared<Message>(tm, protocol::mtMANIFESTS);
1320
1321 manifestListSeq_ = seq;
1322 }
1323
1324 return manifestMessage_;
1325}
1326
1327void
1329 UInt256 const& hash,
1331 std::set<Peer::ID> const& toSkip)
1332{
1333 bool relay = tx.has_value();
1334 if (relay)
1335 {
1336 auto& txn = tx->get();
1337 SerialIter sit(makeSlice(txn.rawtransaction()));
1338 try
1339 {
1340 relay = !isPseudoTx(STTx{sit});
1341 }
1342 catch (std::exception const&)
1343 {
1344 // Could not construct STTx, not relaying
1345 JLOG(journal_.debug()) << "Could not construct STTx: " << hash;
1346 return;
1347 }
1348 }
1349
1350 Overlay::PeerSequence peers = {};
1351 std::size_t total = 0;
1352 std::size_t disabled = 0;
1353 std::size_t enabledInSkip = 0;
1354
1355 if (!relay)
1356 {
1357 if (!app_.config().txReduceRelayEnable)
1358 return;
1359
1360 peers = getActivePeers(toSkip, total, disabled, enabledInSkip);
1361 JLOG(journal_.trace()) << "not relaying tx, total peers " << peers.size();
1362 for (auto const& p : peers)
1363 {
1364 if (p->txReduceRelayEnabled())
1365 p->addTxQueue(hash);
1366 }
1367 return;
1368 }
1369
1370 auto& txn = tx->get();
1371 auto const sm = std::make_shared<Message>(txn, protocol::mtTRANSACTION);
1372 peers = getActivePeers(toSkip, total, disabled, enabledInSkip);
1373 auto const minRelay = app_.config().txReduceRelayMinPeers + disabled;
1374
1375 if (!app_.config().txReduceRelayEnable || total <= minRelay)
1376 {
1377 for (auto const& p : peers)
1378 p->send(sm);
1379 if (app_.config().txReduceRelayEnable || app_.config().txReduceRelayMetrics)
1380 txMetrics_.addMetrics(total, toSkip.size(), 0);
1381 return;
1382 }
1383
1384 // We have more peers than the minimum (disabled + minimum enabled),
1385 // relay to all disabled and some randomly selected enabled that
1386 // do not have the transaction.
1387 auto const enabledTarget = app_.config().txReduceRelayMinPeers +
1388 ((total - minRelay) * app_.config().txRelayPercentage / 100);
1389
1390 txMetrics_.addMetrics(enabledTarget, toSkip.size(), disabled);
1391
1392 if (enabledTarget > enabledInSkip)
1393 std::shuffle(peers.begin(), peers.end(), defaultPrng());
1394
1395 JLOG(journal_.trace()) << "relaying tx, total peers " << peers.size() << " selected "
1396 << enabledTarget << " skip " << toSkip.size() << " disabled "
1397 << disabled;
1398
1399 // count skipped peers with the enabled feature towards the quota
1400 std::uint16_t enabledAndRelayed = enabledInSkip;
1401 for (auto const& p : peers)
1402 {
1403 // always relay to a peer with the disabled feature
1404 if (!p->txReduceRelayEnabled())
1405 {
1406 p->send(sm);
1407 }
1408 else if (enabledAndRelayed < enabledTarget)
1409 {
1410 enabledAndRelayed++;
1411 p->send(sm);
1412 }
1413 else
1414 {
1415 p->addTxQueue(hash);
1416 }
1417 }
1418}
1419
1420//------------------------------------------------------------------------------
1421
1422void
1424{
1426 list_.erase(&child);
1427 if (list_.empty())
1428 cond_.notify_all();
1429}
1430
1431void
1433{
1434 // Calling list_[].second->stop() may cause list_ to be modified
1435 // (OverlayImpl::remove() may be called on this same thread). So
1436 // iterating directly over list_ to call child->stop() could lead to
1437 // undefined behavior.
1438 //
1439 // Therefore we copy all of the weak/shared ptrs out of list_ before we
1440 // start calling stop() on them. That guarantees OverlayImpl::remove()
1441 // won't be called until vector<> children leaves scope.
1443 {
1445 if (!work_)
1446 return;
1447 work_ = std::nullopt;
1448
1449 children.reserve(list_.size());
1450 for (auto const& element : list_)
1451 {
1452 children.emplace_back(element.second.lock());
1453 }
1454 } // lock released
1455
1456 for (auto const& child : children)
1457 {
1458 if (child != nullptr)
1459 child->stop();
1460 }
1461}
1462
1463void
1465{
1466 auto const result = peerFinder_->autoconnect();
1467 for (auto const& addr : result)
1468 connect(addr);
1469}
1470
1471void
1473{
1474 auto const result = peerFinder_->buildEndpointsForPeers();
1475 for (auto const& e : result)
1476 {
1478 {
1480 auto const iter = peers_.find(e.first);
1481 if (iter != peers_.end())
1482 peer = iter->second.lock();
1483 }
1484 if (peer)
1485 peer->sendEndpoints(e.second.begin(), e.second.end());
1486 }
1487}
1488
1489void
1491{
1492 forEach([](auto const& p) {
1493 if (p->txReduceRelayEnabled())
1494 p->sendTxQueue();
1495 });
1496}
1497
1499makeSquelchMessage(PublicKey const& validator, bool squelch, uint32_t squelchDuration)
1500{
1501 protocol::TMSquelch m;
1502 m.set_squelch(squelch);
1503 m.set_validatorpubkey(validator.data(), validator.size());
1504 if (squelch)
1505 m.set_squelchduration(squelchDuration);
1506 return std::make_shared<Message>(m, protocol::mtSQUELCH);
1507}
1508
1509void
1510OverlayImpl::unsquelch(PublicKey const& validator, Peer::ID id) const
1511{
1512 if (auto peer = findPeerByShortID(id); peer)
1513 {
1514 // optimize - multiple message with different
1515 // validator might be sent to the same peer
1516 peer->send(makeSquelchMessage(validator, false, 0));
1517 }
1518}
1519
1520void
1521OverlayImpl::squelch(PublicKey const& validator, Peer::ID id, uint32_t squelchDuration) const
1522{
1523 if (auto peer = findPeerByShortID(id); peer)
1524 {
1525 peer->send(makeSquelchMessage(validator, true, squelchDuration));
1526 }
1527}
1528
1529void
1531 UInt256 const& key,
1532 PublicKey const& validator,
1533 std::set<Peer::ID>&& peers,
1534 protocol::MessageType type)
1535{
1536 if (!slots_.baseSquelchReady())
1537 return;
1538
1539 if (!strand_.running_in_this_thread())
1540 {
1541 post(
1542 strand_,
1543 // Must capture copies of reference parameters (i.e. key, validator)
1544 [this, key = key, validator = validator, peers = std::move(peers), type]() mutable {
1545 updateSlotAndSquelch(key, validator, std::move(peers), type);
1546 });
1547
1548 return;
1549 }
1550
1551 for (auto id : peers)
1552 {
1553 slots_.updateSlotAndSquelch(key, validator, id, type, [&]() {
1555 });
1556 }
1557}
1558
1559void
1561 UInt256 const& key,
1562 PublicKey const& validator,
1563 Peer::ID peer,
1564 protocol::MessageType type)
1565{
1566 if (!slots_.baseSquelchReady())
1567 return;
1568
1569 if (!strand_.running_in_this_thread())
1570 {
1571 {
1572 post(
1573 strand_,
1574 // Must capture copies of reference parameters (i.e. key, validator)
1575 [this, key = key, validator = validator, peer, type]() {
1576 updateSlotAndSquelch(key, validator, peer, type);
1577 });
1578 }
1579 return;
1580 }
1581
1582 slots_.updateSlotAndSquelch(key, validator, peer, type, [&]() {
1584 });
1585}
1586
1587void
1589{
1590 if (!strand_.running_in_this_thread())
1591 {
1592 post(strand_, [this, id] { deletePeer(id); });
1593 return;
1594 }
1595
1596 slots_.deletePeer(id, true);
1597}
1598
1599void
1601{
1602 if (!strand_.running_in_this_thread())
1603 {
1604 post(strand_, [this] { deleteIdlePeers(); });
1605 return;
1606 }
1607
1608 slots_.deleteIdlePeers();
1609}
1610
1611//------------------------------------------------------------------------------
1612
1615{
1616 Overlay::Setup setup;
1617
1618 {
1619 auto const& section = config.section(Sections::kOverlay);
1620 setup.context = makeSslContext("");
1621
1622 set(setup.ipLimit, "ip_limit", section);
1623 if (setup.ipLimit < 0)
1624 Throw<std::runtime_error>("Configured IP limit is invalid");
1625
1626 std::string ip;
1627 set(ip, "public_ip", section);
1628 if (!ip.empty())
1629 {
1630 boost::system::error_code ec;
1631 setup.publicIp = boost::asio::ip::make_address(ip, ec);
1632 if (ec || !beast::ip::isPublic(setup.publicIp))
1633 Throw<std::runtime_error>("Configured public IP is invalid");
1634 }
1635
1636 set(setup.verifyEndpoints, true, "verify_endpoints", section);
1637 if (!setup.verifyEndpoints)
1638 {
1639 JLOG(j.warn()) << "Endpoint verification is disabled. This is a "
1640 "security risk and should only be used for "
1641 "testing.";
1642 }
1643 }
1644
1645 {
1646 auto const& section = config.section(Sections::kCrawl);
1647 auto const& values = section.values();
1648
1649 if (values.size() > 1)
1650 {
1651 Throw<std::runtime_error>("Configured [crawl] section is invalid, too many values");
1652 }
1653
1654 bool crawlEnabled = true;
1655
1656 // Only allow "0|1" as a value
1657 if (values.size() == 1)
1658 {
1659 try
1660 {
1661 crawlEnabled = boost::lexical_cast<bool>(values.front());
1662 }
1663 catch (boost::bad_lexical_cast const&)
1664 {
1666 "Configured [crawl] section has invalid value: {}", values.front()));
1667 }
1668 }
1669
1670 if (crawlEnabled)
1671 {
1672 if (get<bool>(section, Keys::kOverlay, true))
1673 {
1675 }
1676 if (get<bool>(section, Keys::kServer, true))
1677 {
1679 }
1680 if (get<bool>(section, Keys::kCounts, false))
1681 {
1683 }
1684 if (get<bool>(section, Keys::kUnl, true))
1685 {
1687 }
1688 }
1689 }
1690 {
1691 auto const& section = config.section(Sections::kVl);
1692
1693 set(setup.vlEnabled, "enabled", section);
1694 }
1695
1696 try
1697 {
1698 auto id = config.legacy(Sections::kNetworkId);
1699
1700 if (!id.empty())
1701 {
1702 if (id == "main")
1703 id = "0";
1704
1705 if (id == "testnet")
1706 id = "1";
1707
1708 if (id == "devnet")
1709 id = "2";
1710
1712 }
1713 }
1714 catch (...)
1715 {
1717 "Configured [network_id] section is invalid: must be a number "
1718 "or one of the strings 'main', 'testnet' or 'devnet'.");
1719 }
1720
1721 return setup;
1722}
1723
1726 Application& app,
1727 Overlay::Setup const& setup,
1728 ServerHandler& serverHandler,
1729 resource::Manager& resourceManager,
1730 Resolver& resolver,
1731 boost::asio::io_context& ioContext,
1732 BasicConfig const& config,
1733 beast::insight::Collector::Ptr const& collector)
1734{
1736 app, setup, serverHandler, resourceManager, resolver, ioContext, config, collector);
1737}
1738
1739} // namespace xrpl
T begin(T... args)
A generic endpoint for log messages.
Definition Journal.h:44
Stream debug() const
Definition Journal.h:344
Stream warn() const
Definition Journal.h:356
std::string const & name() const
Returns the name of this source.
void add(Source &source)
Add a child source.
Wraps a Journal::Sink to prefix its output with a string.
Definition WrappedSink.h:19
std::shared_ptr< Collector > Ptr
Definition Collector.h:29
A version-independent IP address and port combination.
Definition IPEndpoint.h:24
Represents a JSON value.
Definition json_value.h:117
Value removeMember(char const *key)
Remove and return the named member.
Value & append(Value const &value)
Append value to array at the end.
bool isMember(char const *key) const
Return true if the object has a member named key.
Holds unparsed configuration information.
void legacy(std::string const &section, std::string value)
Set a value that is not a key/value pair.
Section & section(std::string const &name)
Returns the section with the given name.
Child(OverlayImpl &overlay)
static std::shared_ptr< Writer > makeErrorResponse(std::shared_ptr< peer_finder::Slot > const &slot, HttpRequestType const &request, AddressType remoteAddress, std::string const &msg)
boost::system::error_code ErrorCode
Definition OverlayImpl.h:90
boost::asio::io_context & ioContext_
resource::Manager & resourceManager()
void squelch(PublicKey const &validator, Peer::ID const id, std::uint32_t squelchDuration) const override
Squelch handler.
std::set< Peer::ID > relay(protocol::TMProposeSet const &m, UInt256 const &uid, PublicKey const &validator) override
Relay a proposal.
json::Value getOverlayInfo() const
Returns information about peers on the overlay network.
bool processRequest(HttpRequestType const &req, Handoff &handoff)
Handles non-peer protocol requests.
bool processValidatorList(HttpRequestType const &req, Handoff &handoff)
Handles validator list requests.
Resolver & resolver_
void addActive(std::shared_ptr< PeerImp > const &peer)
std::atomic< Peer::ID > nextId_
void broadcast(protocol::TMProposeSet const &m) override
Broadcast a proposal.
void activate(std::shared_ptr< PeerImp > const &peer)
Called when a peer has connected successfully This is called after the peer handshake has been comple...
std::optional< boost::asio::executor_work_guard< boost::asio::io_context::executor_type > > work_
void connect(beast::ip::Endpoint const &remoteEndpoint) override
Establish a peer connection to the specified endpoint.
peer_finder::Manager & peerFinder()
void remove(std::shared_ptr< peer_finder::Slot > const &slot)
void onPeerDeactivate(Peer::ID id)
void stop() override
std::size_t size() const override
The number of active peers on the network Active peers are only those peers that have completed the h...
static bool isUpgrade(boost::beast::http::header< true, Fields > const &req)
ServerHandler & serverHandler_
void onManifests(std::shared_ptr< protocol::TMManifests > const &m, std::shared_ptr< PeerImp > const &from)
std::shared_ptr< Writer > makeRedirectResponse(std::shared_ptr< peer_finder::Slot > const &slot, HttpRequestType const &request, AddressType remoteAddress)
peer_finder::StoreSqdb store_
reduce_relay::Slots< UptimeClock > slots_
void deleteIdlePeers()
Check if peers stopped relaying messages and if slots stopped receiving messages from the validator.
void unsquelch(PublicKey const &validator, Peer::ID id) const override
Unsquelch handler.
bool processHealth(HttpRequestType const &req, Handoff &handoff)
Handles health requests.
void reportInboundTraffic(TrafficCount::Category cat, int bytes)
std::shared_ptr< Message > manifestMessage_
void sendTxQueue() const
Send once a second transactions' hashes aggregated by peers.
std::optional< std::uint32_t > manifestListSeq_
std::unique_ptr< peer_finder::Manager > peerFinder_
void onWrite(beast::PropertyStream::Map &stream) override
Subclass override.
resource::Manager & resourceManager_
HashMap< std::shared_ptr< peer_finder::Slot >, std::weak_ptr< PeerImp > > peers_
Application & app_
json::Value getServerCounts()
Returns information about the local server's performance counters.
std::recursive_mutex mutex_
bool processCrawl(HttpRequestType const &req, Handoff &handoff)
Handles crawl requests.
beast::Journal const journal_
json::Value json() override
Return diagnostics on the status of all peers.
void forEach(UnaryFunc &&f) const
std::shared_ptr< Peer > findPeerByShortID(Peer::ID const &id) const override
Returns the peer with the matching short id, or null.
std::mutex manifestLock_
Handoff onHandoff(std::unique_ptr< StreamType > &&bundle, HttpRequestType &&request, EndpointType remoteEndpoint) override
Conditionally accept an incoming HTTP request.
boost::asio::strand< boost::asio::io_context::executor_type > strand_
static std::string makePrefix(std::uint32_t id)
Setup const & setup() const
HashMap< Peer::ID, std::weak_ptr< PeerImp > > ids_
boost::asio::ip::address AddressType
Definition OverlayImpl.h:88
static bool isPeerUpgrade(HttpRequestType const &request)
metrics::TxMetrics txMetrics_
boost::container::flat_map< Child *, std::weak_ptr< Child > > list_
int limit() override
Returns the maximum number of peers we are configured to allow.
std::condition_variable_any cond_
OverlayImpl(Application &app, Setup setup, ServerHandler &serverHandler, resource::Manager &resourceManager, Resolver &resolver, boost::asio::io_context &ioContext, BasicConfig const &config, beast::insight::Collector::Ptr const &collector)
json::Value getUnlInfo()
Returns information about the local server's UNL.
std::shared_ptr< Message > getManifestsMessage()
std::shared_ptr< Peer > findPeerByPublicKey(PublicKey const &pubKey) override
Returns the peer with the matching public key, or null.
TrafficCount traffic_
boost::asio::ip::tcp::endpoint EndpointType
Definition OverlayImpl.h:89
void deletePeer(Peer::ID id)
Called when the peer is deleted.
void checkTracking(std::uint32_t) override
Calls the checkTracking function on each peer.
json::Value getServerInfo()
Returns information about the local server.
void reportOutboundTraffic(TrafficCount::Category cat, int bytes)
void start() override
PeerSequence getActivePeers() const override
Returns a sequence representing the current list of peers.
void updateSlotAndSquelch(UInt256 const &key, PublicKey const &validator, std::set< Peer::ID > &&peers, protocol::MessageType type)
Updates message count for validator/peer.
std::vector< std::shared_ptr< Peer > > PeerSequence
Definition Overlay.h:66
std::uint32_t ID
Uniquely identifies a peer.
A public key.
Definition PublicKey.h:53
std::vector< std::string > const & values() const
Returns all the values in the section.
Definition BasicConfig.h:68
virtual std::pair< std::shared_ptr< Slot >, Result > newOutboundSlot(beast::ip::Endpoint const &remoteEndpoint)=0
Create a new outbound slot with the specified remote endpoint.
Tracks load and resource consumption.
virtual Consumer newOutboundEndpoint(beast::ip::Endpoint const &address)=0
Create a new endpoint keyed by outbound IP address and port.
T contains(T... args)
T duration_cast(T... args)
T emplace_back(T... args)
T emplace(T... args)
T empty(T... args)
T end(T... args)
T find_if(T... args)
T format(T... args)
T get(T... args)
T insert(T... args)
T lock(T... args)
T make_shared(T... args)
T make_tuple(T... args)
T make_unique(T... args)
T max(T... args)
T min(T... args)
bool isPublic(Address const &addr)
Returns true if the address is a public routable address.
Definition IPAddress.h:71
bool isKeepAlive(boost::beast::http::message< IsRequest, Body, Fields > const &m)
Definition rfc2616.h:367
Result splitCommas(FwdIt first, FwdIt last)
Definition rfc2616.h:183
constexpr Out lexicalCastThrow(In in)
Convert from one type to another, throw on error.
JSON (JavaScript Object Notation).
Definition json_errors.h:5
@ Array
array value (ordered list)
Definition json_value.h:28
@ Object
object value (collection of name/value pairs).
Definition json_value.h:29
@ Null
'null' value
Definition json_value.h:22
STL namespace.
std::string const & getFullVersionString()
Full server version string.
Definition BuildInfo.cpp:77
static constexpr auto kDisabled
static constexpr auto kOverlay
static constexpr auto kUnl
static constexpr auto kServerCounts
static constexpr auto kServerInfo
Config makeConfig(xrpl::Config const &cfg, std::uint16_t port, bool validationPublicKey, int ipLimit, bool verifyEndpoints)
Charge const kFeeMalformedRequest
Schedule of fees charged for imposing load on the server.
static constexpr auto kCheckIdlePeers
How often we check for idle peers (seconds).
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,...
beast::XorShiftEngine & defaultPrng()
Return the default random engine.
Stopwatch & stopwatch()
Returns an instance of a wall clock.
Definition chrono.h:101
T get(Section const &section, std::string const &name, T const &defaultValue=T{})
Retrieve a key/value pair from a section.
std::string strHex(FwdIt begin, FwdIt end)
Definition strHex.h:13
std::optional< ProtocolVersion > negotiateProtocolVersion(std::vector< ProtocolVersion > const &versions)
Given a list of supported protocol versions, choose the one we prefer.
std::string to_string(BaseUInt< Bits, Tag > const &a)
Definition base_uint.h:657
void addValidatorManifest(soci::session &session, std::string const &serialized)
addValidatorManifest Saves the manifest of a validator to the database.
Definition Wallet.cpp:137
std::optional< Manifest > deserializeManifest(Slice s, beast::Journal journal)
Constructs Manifest from serialized string.
@ Uncapped
Bypasses the cap (listed/trusted or config manifests).
Definition Manifest.h:365
@ Capped
Subject to the untrusted cap (unlisted peer gossip).
Definition Manifest.h:364
Slice makeSlice(std::array< T, N > const &a)
Definition Slice.h:228
BaseUInt< 256 > UInt256
Definition base_uint.h:580
boost::beast::http::request< boost::beast::http::dynamic_body > HttpRequestType
Definition Handoff.h:12
std::shared_ptr< boost::asio::ssl::context > makeSslContext(std::string const &cipherList)
Create a self-signed SSL context that allows anonymous Diffie Hellman.
std::string base64Encode(std::uint8_t const *data, std::size_t len)
std::optional< UInt256 > makeSharedValue(StreamType &ssl, beast::Journal journal)
Computes a shared value based on the SSL connection state.
std::unique_ptr< Overlay > makeOverlay(Application &app, Overlay::Setup const &setup, ServerHandler &serverHandler, resource::Manager &resourceManager, Resolver &resolver, boost::asio::io_context &ioContext, BasicConfig const &config, beast::insight::Collector::Ptr const &collector)
Creates the implementation of Overlay.
constexpr Number squelch(Number const &x, Number const &limit) noexcept
Definition Number.h:907
std::shared_ptr< Message > makeSquelchMessage(PublicKey const &validator, bool squelch, uint32_t squelchDuration)
constexpr std::size_t untrustedManifestCount(std::optional< std::size_t > const &configured)
Number of untrusted manifests to store in cache and allowed in one Manifest message.
Definition Manifest.h:246
Overlay::Setup setupOverlay(BasicConfig const &config, beast::Journal j)
json::Value getCountsJson(Application &app, int minObjectCount)
Definition GetCounts.cpp:46
std::vector< ProtocolVersion > parseProtocolVersions(std::string_view value)
Parse a set of protocol versions.
bool isPseudoTx(STObject const &tx)
Check whether a transaction is a pseudo-transaction.
Definition STTx.cpp:889
XRPL_NO_SANITIZE_ADDRESS void Throw(Args &&... args)
Definition contract.h:52
PublicKey verifyHandshake(boost::beast::http::fields const &headers, xrpl::UInt256 const &sharedValue, std::optional< std::uint32_t > networkID, beast::ip::Address publicIp, beast::ip::Address remote, Application &app)
Validate header fields necessary for upgrading the link to the peer protocol.
@ Accepted
Manifest is valid.
Definition Manifest.h:321
T piecewise_construct
T push_back(T... args)
T shuffle(T... args)
T reserve(T... args)
T setfill(T... args)
T setw(T... args)
T size(T... args)
T str(T... args)
static boost::asio::ip::tcp::endpoint toAsioEndpoint(ip::Endpoint const &address)
static ip::Endpoint fromAsio(boost::asio::ip::address const &address)
Used to indicate the result of a server connection handoff.
Definition Handoff.h:20
bool keepAlive
Definition Handoff.h:26
std::shared_ptr< Writer > response
Definition Handoff.h:29
static constexpr auto kUnl
Definition Constants.h:176
static constexpr auto kCounts
Definition Constants.h:103
static constexpr auto kOverlay
Definition Constants.h:140
static constexpr auto kServer
Definition Constants.h:157
void onTimer(ErrorCode ec)
boost::asio::basic_waitable_timer< ClockType > timer
Definition OverlayImpl.h:94
Timer(OverlayImpl &overlay)
std::uint32_t crawlOptions
Definition Overlay.h:60
std::optional< std::uint32_t > networkID
Definition Overlay.h:61
beast::ip::Address publicIp
Definition Overlay.h:58
std::shared_ptr< boost::asio::ssl::context > context
Definition Overlay.h:57
static constexpr auto kOverlay
Definition Constants.h:35
static constexpr auto kVl
Definition Constants.h:77
static constexpr auto kCrawl
Definition Constants.h:12
static constexpr auto kNetworkId
Definition Constants.h:30
PeerFinder configuration settings.
T substr(T... args)
T to_string(T... args)
T what(T... args)