xrpld
Loading...
Searching...
No Matches
NetworkOPs.cpp
1#include <xrpl/server/NetworkOPs.h>
2
3#include <xrpld/app/consensus/RCLConsensus.h>
4#include <xrpld/app/consensus/RCLCxPeerPos.h>
5#include <xrpld/app/consensus/RCLValidations.h>
6#include <xrpld/app/ledger/AcceptedLedger.h>
7#include <xrpld/app/ledger/InboundLedger.h>
8#include <xrpld/app/ledger/InboundLedgers.h>
9#include <xrpld/app/ledger/LedgerMaster.h>
10#include <xrpld/app/ledger/LedgerToJson.h>
11#include <xrpld/app/ledger/LocalTxs.h>
12#include <xrpld/app/ledger/OpenLedger.h>
13#include <xrpld/app/ledger/TransactionMaster.h>
14#include <xrpld/app/main/LoadManager.h>
15#include <xrpld/app/main/Tuning.h>
16#include <xrpld/app/misc/DeliverMax.h>
17#include <xrpld/app/misc/FeeVote.h>
18#include <xrpld/app/misc/Transaction.h>
19#include <xrpld/app/misc/TxQ.h>
20#include <xrpld/app/misc/ValidatorKeys.h>
21#include <xrpld/app/misc/ValidatorList.h>
22#include <xrpld/app/misc/make_NetworkOPs.h>
23#include <xrpld/app/rdb/backend/SQLiteDatabase.h>
24#include <xrpld/core/Config.h>
25#include <xrpld/overlay/Cluster.h>
26#include <xrpld/overlay/ClusterNode.h>
27#include <xrpld/overlay/Overlay.h>
28#include <xrpld/overlay/predicates.h>
29#include <xrpld/rpc/BookChanges.h>
30#include <xrpld/rpc/CTID.h>
31#include <xrpld/rpc/ServerHandler.h>
32#include <xrpld/rpc/detail/SyntheticFields.h>
33
34#include <xrpl/basics/Log.h>
35#include <xrpl/basics/ToString.h>
36#include <xrpl/basics/UnorderedContainers.h>
37#include <xrpl/basics/UptimeClock.h>
38#include <xrpl/basics/base_uint.h>
39#include <xrpl/basics/chrono.h>
40#include <xrpl/basics/contract.h>
41#include <xrpl/basics/mulDiv.h>
42#include <xrpl/basics/safe_cast.h>
43#include <xrpl/basics/scope.h>
44#include <xrpl/basics/strHex.h>
45#include <xrpl/beast/clock/abstract_clock.h>
46#include <xrpl/beast/insight/Collector.h>
47#include <xrpl/beast/insight/Gauge.h>
48#include <xrpl/beast/insight/Hook.h>
49#include <xrpl/beast/net/IPEndpoint.h>
50#include <xrpl/beast/utility/Zero.h>
51#include <xrpl/beast/utility/instrumentation.h>
52#include <xrpl/beast/utility/rngfill.h>
53#include <xrpl/config/Constants.h>
54#include <xrpl/consensus/ConsensusParms.h>
55#include <xrpl/consensus/ConsensusTypes.h>
56#include <xrpl/core/ClosureCounter.h>
57#include <xrpl/core/HashRouter.h>
58#include <xrpl/core/Job.h>
59#include <xrpl/core/NetworkIDService.h>
60#include <xrpl/core/PerfLog.h>
61#include <xrpl/core/ServiceRegistry.h>
62#include <xrpl/crypto/RFC1751.h>
63#include <xrpl/crypto/csprng.h>
64#include <xrpl/git/Git.h>
65#include <xrpl/json/json_forwards.h>
66#include <xrpl/json/json_value.h>
67#include <xrpl/json/json_writer.h>
68#include <xrpl/ledger/AcceptedLedgerTx.h>
69#include <xrpl/ledger/AmendmentTable.h>
70#include <xrpl/ledger/ApplyView.h>
71#include <xrpl/ledger/CanonicalTXSet.h>
72#include <xrpl/ledger/Ledger.h>
73#include <xrpl/ledger/OpenView.h>
74#include <xrpl/ledger/OrderBookDB.h>
75#include <xrpl/ledger/ReadView.h>
76#include <xrpl/ledger/helpers/DirectoryHelpers.h>
77#include <xrpl/ledger/helpers/MPTokenHelpers.h>
78#include <xrpl/ledger/helpers/TokenHelpers.h>
79#include <xrpl/protocol/AccountID.h>
80#include <xrpl/protocol/AmountConversions.h>
81#include <xrpl/protocol/ApiVersion.h>
82#include <xrpl/protocol/Book.h>
83#include <xrpl/protocol/BuildInfo.h>
84#include <xrpl/protocol/ErrorCodes.h>
85#include <xrpl/protocol/Feature.h>
86#include <xrpl/protocol/Fees.h>
87#include <xrpl/protocol/Indexes.h>
88#include <xrpl/protocol/Issue.h>
89#include <xrpl/protocol/KeyType.h>
90#include <xrpl/protocol/LedgerFormats.h>
91#include <xrpl/protocol/MPTAmount.h>
92#include <xrpl/protocol/MPTIssue.h>
93#include <xrpl/protocol/MultiApiJson.h>
94#include <xrpl/protocol/Protocol.h>
95#include <xrpl/protocol/PublicKey.h>
96#include <xrpl/protocol/RPCErr.h>
97#include <xrpl/protocol/Rate.h>
98#include <xrpl/protocol/SField.h>
99#include <xrpl/protocol/STAmount.h>
100#include <xrpl/protocol/STTx.h>
101#include <xrpl/protocol/SecretKey.h>
102#include <xrpl/protocol/Seed.h>
103#include <xrpl/protocol/Serializer.h>
104#include <xrpl/protocol/TER.h>
105#include <xrpl/protocol/TxFlags.h>
106#include <xrpl/protocol/TxFormats.h>
107#include <xrpl/protocol/UintTypes.h>
108#include <xrpl/protocol/Units.h>
109#include <xrpl/protocol/XRPAmount.h>
110#include <xrpl/protocol/jss.h>
111#include <xrpl/protocol/tokens.h>
112#include <xrpl/rdb/RelationalDatabase.h>
113#include <xrpl/resource/Fees.h>
114#include <xrpl/resource/Gossip.h>
115#include <xrpl/resource/ResourceManager.h>
116#include <xrpl/server/InfoSub.h>
117#include <xrpl/server/LoadFeeTrack.h>
118#include <xrpl/server/Manifest.h>
119#include <xrpl/shamap/SHAMap.h>
120#include <xrpl/tx/apply.h>
121
122#include <boost/asio/error.hpp>
123#include <boost/asio/io_context.hpp>
124#include <boost/asio/ip/host_name.hpp>
125#include <boost/asio/steady_timer.hpp>
126#include <boost/system/detail/errc.hpp>
127#include <boost/system/detail/error_code.hpp>
128#include <boost/system/system_error.hpp>
129
130#include <xrpl.pb.h>
131
132#include <algorithm>
133#include <array>
134#include <atomic>
135#include <chrono>
136#include <climits>
137#include <condition_variable>
138#include <cstddef>
139#include <cstdint>
140#include <cstdlib>
141#include <exception>
142#include <functional>
143#include <iterator>
144#include <limits>
145#include <memory>
146#include <mutex>
147#include <optional>
148#include <set>
149#include <sstream>
150#include <stdexcept>
151#include <string>
152#include <string_view>
153#include <type_traits>
154#include <unordered_map>
155#include <utility>
156#include <vector>
157
158namespace xrpl {
159
167class NetworkOPsImp final : public NetworkOPs
168{
172
174 {
175 public:
177 bool const admin;
178 bool const local;
180 bool applied = false;
182
184 : transaction(std::move(t)), admin(a), local(l), failType(f)
185 {
186 XRPL_ASSERT(
188 "xrpl::NetworkOPsImp::TransactionStatus::TransactionStatus : "
189 "valid inputs");
190 }
191 };
192
196 enum class DispatchState : unsigned char {
200 };
201
203
219 {
227
231 std::chrono::steady_clock::time_point start_ = std::chrono::steady_clock::now();
232 std::chrono::steady_clock::time_point const processStart_ = start_;
235
236 public:
238 {
239 counters_[static_cast<std::size_t>(OperatingMode::DISCONNECTED)].transitions = 1;
240 }
241
248 void
250
256 void
257 json(json::Value& obj) const;
258
266
267 CounterData
269 {
270 std::scoped_lock const lock(mutex_);
271 return {
272 .counters = counters_,
273 .mode = mode_,
274 .start = start_,
275 .initialSyncUs = initialSyncUs_};
276 }
277 };
278
283 {
284 ServerFeeSummary() = default;
285
287 XRPAmount fee,
288 TxQ::Metrics escalationMetrics, // trivially copyable
289 LoadFeeTrack const& loadFeeTrack);
290 bool
291 operator!=(ServerFeeSummary const& b) const;
292
293 bool
295 {
296 return !(*this != b);
297 }
298
303 };
304
305public:
307 ServiceRegistry& registry,
309 bool standalone,
310 std::size_t minPeerCount,
311 bool startValid,
312 JobQueue& jobQueue,
313 LedgerMaster& ledgerMaster,
314 ValidatorKeys const& validatorKeys,
315 boost::asio::io_context& ioCtx,
317 beast::insight::Collector::Ptr const& collector)
318 : registry_(registry)
322 , heartbeatTimer_(ioCtx)
323 , clusterTimer_(ioCtx)
325 , consensus_(
326 registry_.get().getApp(),
328 setupFeeVote(registry_.get().getApp().config().section(Sections::kVoting)),
329 registry_.get().getJournal("FeeVote")),
330 ledgerMaster,
331 *localTX_,
332 registry.getInboundTransactions(),
333 beast::getAbstractClock<std::chrono::steady_clock>(),
334 validatorKeys,
335 registry_.get().getJournal("LedgerConsensus"))
336 , validatorPK_(
337 validatorKeys.keys ? validatorKeys.keys->publicKey : decltype(validatorPK_){})
339 validatorKeys.keys ? validatorKeys.keys->masterPublicKey
340 : decltype(validatorMasterPK_){})
341 , ledgerMaster_(ledgerMaster)
342 , jobQueue_(jobQueue)
343 , standalone_(standalone)
344 , minPeerCount_(startValid ? 0 : minPeerCount)
345 , stats_([this] { collectMetrics(); }, collector)
346 {
347 }
348
349 ~NetworkOPsImp() override
350 {
351 // This clear() is necessary to ensure the shared_ptrs in this map get
352 // destroyed NOW because the objects in this map invoke methods on this
353 // class when they are destroyed
354 rpcSubMap_.clear();
355 }
356
357public:
359 getOperatingMode() const override;
360
362 strOperatingMode(OperatingMode const mode, bool const admin) const override;
363
365 strOperatingMode(bool const admin = false) const override;
366
367 //
368 // Transaction operations.
369 //
370
371 // Must complete immediately.
372 void
374
375 void
377 std::shared_ptr<Transaction>& transaction,
378 bool bUnlimited,
379 bool bLocal,
380 FailHard failType) override;
381
382 void
384
393 void
394 doTransactionSync(std::shared_ptr<Transaction> transaction, bool bUnlimited, FailHard failType);
395
405 void
408 bool bUnlimited,
409 FailHard failtype);
410
411private:
412 bool
414
415 void
418 std::function<bool(std::unique_lock<std::mutex> const&)> retryCallback);
419
420public:
424 void
426
432 void
434
435 //
436 // Owner functions.
437 //
438
440 getOwnerInfo(std::shared_ptr<ReadView const> lpLedger, AccountID const& account) override;
441
442 //
443 // Book functions.
444 //
445
446 void
449 Book const&,
450 AccountID const& uTakerID,
451 bool const bProof,
452 unsigned int iLimit,
453 json::Value const& jvMarker,
454 json::Value& jvResult) override;
455
456 // Ledger proposal/close functions.
457 bool
459
460 bool
461 recvValidation(std::shared_ptr<STValidation> const& val, std::string const& source) override;
462
463 void
464 mapComplete(std::shared_ptr<SHAMap> const& map, bool fromAcquire) override;
465
466 // Network state machine.
467
468 // Used for the "jump" case.
469private:
470 void
472 bool
474
475public:
476 bool
477 beginConsensus(UInt256 const& networkClosed, std::unique_ptr<std::stringstream> const& clog)
478 override;
479 void
481 void
482 setStandAlone() override;
483
488 void
489 setStateTimer() override;
490
491 void
492 setNeedNetworkLedger() override;
493 void
494 clearNeedNetworkLedger() override;
495 bool
496 isNeedNetworkLedger() override;
497 bool
498 isFull() override;
499
500 void
501 setMode(OperatingMode om) override;
502
503 bool
504 isBlocked() override;
505 bool
506 isAmendmentBlocked() override;
507 void
508 setAmendmentBlocked() override;
509 bool
510 isAmendmentWarned() override;
511 void
512 setAmendmentWarned() override;
513 void
514 clearAmendmentWarned() override;
515 bool
516 isUNLBlocked() override;
517 void
518 setUNLBlocked() override;
519 void
520 clearUNLBlocked() override;
521 void
522 consensusViewChange() override;
523
525 getConsensusInfo() override;
527 getServerInfo(bool human, bool admin, bool counters) override;
528 void
529 clearLedgerFetch() override;
531 getLedgerFetchInfo() override;
534 void
535 reportFeeChange() override;
536 void
538
539 void
540 updateLocalTx(ReadView const& view) override;
542 getLocalTxCount() override;
544 getBookSubscribersCount() override;
545
546 //
547 // Monitoring: publisher side.
548 //
549 void
550 pubLedger(std::shared_ptr<ReadView const> const& lpAccepted) override;
551 void
554 std::shared_ptr<STTx const> const& transaction,
555 TER result) override;
556 void
557 pubValidation(std::shared_ptr<STValidation> const& val) override;
558
559 //--------------------------------------------------------------------------
560 //
561 // InfoSub::Source.
562 //
563 void
564 subAccount(InfoSub::Ref ispListener, HashSet<AccountID> const& vnaAccountIDs, bool rt) override;
565 void
566 unsubAccount(InfoSub::Ref ispListener, HashSet<AccountID> const& vnaAccountIDs, bool rt)
567 override;
568
569 // Just remove the subscription from the tracking
570 // not from the InfoSub. Needed for InfoSub destruction
571 void
572 unsubAccountInternal(std::uint64_t seq, HashSet<AccountID> const& vnaAccountIDs, bool rt)
573 override;
574
575 void
576 subMPT(InfoSub::Ref ispListener, HashSet<MPTID> const& mptIDs) override;
577 void
578 unsubMPT(InfoSub::Ref ispListener, HashSet<MPTID> const& mptIDs) override;
579 void
580 unsubMPTInternal(std::uint64_t seq, MPTID const& mptID) override;
581
583 subAccountHistory(InfoSub::Ref ispListener, AccountID const& account) override;
584 void
585 unsubAccountHistory(InfoSub::Ref ispListener, AccountID const& account, bool historyOnly)
586 override;
587
588 void
589 unsubAccountHistoryInternal(std::uint64_t seq, AccountID const& account, bool historyOnly)
590 override;
591
592 void
594 std::uint64_t seq,
595 HashSet<AccountID> rtAccounts,
596 HashSet<AccountID> normalAccounts,
597 HashSet<AccountID> historyAccounts) override;
598
599 bool
600 subLedger(InfoSub::Ref ispListener, json::Value& jvResult) override;
601 bool
602 unsubLedger(std::uint64_t uListener) override;
603
604 bool
605 subBookChanges(InfoSub::Ref ispListener) override;
606 bool
607 unsubBookChanges(std::uint64_t uListener) override;
608
609 bool
610 subServer(InfoSub::Ref ispListener, json::Value& jvResult, bool admin) override;
611 bool
612 unsubServer(std::uint64_t uListener) override;
613
614 bool
615 subBook(InfoSub::Ref ispListener, Book const&) override;
616 bool
617 unsubBook(InfoSub::Ref ispListener, Book const&) override;
618 bool
619 unsubBookInternal(std::uint64_t uListener, Book const&) override;
620
621 bool
622 subManifests(InfoSub::Ref ispListener) override;
623 bool
624 unsubManifests(std::uint64_t uListener) override;
625 void
626 pubManifest(Manifest const&) override;
627
628 bool
629 subTransactions(InfoSub::Ref ispListener) override;
630 bool
631 unsubTransactions(std::uint64_t uListener) override;
632
633 bool
634 subRTTransactions(InfoSub::Ref ispListener) override;
635 bool
636 unsubRTTransactions(std::uint64_t uListener) override;
637
638 bool
639 subValidations(InfoSub::Ref ispListener) override;
640 bool
641 unsubValidations(std::uint64_t uListener) override;
642
643 bool
644 subPeerStatus(InfoSub::Ref ispListener) override;
645 bool
646 unsubPeerStatus(std::uint64_t uListener) override;
647 void
648 pubPeerStatus(std::function<json::Value(void)> const&) override;
649
650 bool
651 subConsensus(InfoSub::Ref ispListener) override;
652 bool
653 unsubConsensus(std::uint64_t uListener) override;
654
656 findRpcSub(std::string const& strUrl) override;
658 addRpcSub(std::string const& strUrl, InfoSub::Ref) override;
659 bool
660 tryRemoveRpcSub(std::string const& strUrl) override;
661
674 findRpcSubLocked(std::string const& strUrl);
675
676 beast::Journal const&
677 journal() const override
678 {
679 return journal_;
680 }
681
682 void
683 stop() override
684 {
685 {
686 try
687 {
688 heartbeatTimer_.cancel();
689 }
690 catch (boost::system::system_error const& e)
691 {
692 JLOG(journal_.error()) << "NetworkOPs: heartbeatTimer cancel error: " << e.what();
693 }
694
695 try
696 {
697 clusterTimer_.cancel();
698 }
699 catch (boost::system::system_error const& e)
700 {
701 JLOG(journal_.error()) << "NetworkOPs: clusterTimer cancel error: " << e.what();
702 }
703
704 try
705 {
706 accountHistoryTxTimer_.cancel();
707 }
708 catch (boost::system::system_error const& e)
709 {
710 JLOG(journal_.error())
711 << "NetworkOPs: accountHistoryTxTimer cancel error: " << e.what();
712 }
713 }
714 // Make sure that any waitHandlers pending in our timers are done.
715 using namespace std::chrono_literals;
716 waitHandlerCounter_.join("NetworkOPs", 1s, journal_);
717 }
718
719 void
720 stateAccounting(json::Value& obj) override;
721
722private:
723 void
724 setTimer(
725 boost::asio::steady_timer& timer,
726 std::chrono::milliseconds const& expiryTime,
727 std::function<void()> onExpire,
728 std::function<void()> onError);
729 void
731 void
733 void
735 void
737
739 transJson(
740 std::shared_ptr<STTx const> const& transaction,
741 TER result,
742 bool validated,
745
746 void
749 AcceptedLedgerTx const& transaction,
750 bool last);
751
752 void
755 AcceptedLedgerTx const& transaction,
756 bool last);
757
781 void
782 pubBookTransaction(AcceptedLedgerTx const& transaction, MultiApiJson const& jvObj);
783
784 void
787 std::shared_ptr<STTx const> const& transaction,
788 TER result);
789
794 void
796 std::shared_ptr<ReadView const> const& lpAccepted,
797 std::shared_ptr<AcceptedLedger const> const& alpAccepted);
798
804 void
806
807 void
808 pubServer();
809 void
811 void
812 pubMPTTransaction(AcceptedLedgerTx const& transaction, MultiApiJson const& jvObj);
813
815 getHostId(bool forAdmin);
816
817private:
822
823 /*
824 * With a validated ledger to separate history and future, the node
825 * streams historical txns with negative indexes starting from -1,
826 * and streams future txns starting from index 0.
827 * The SubAccountHistoryIndex struct maintains these indexes.
828 * It also has a flag stopHistorical_ for stopping streaming
829 * the historical txns.
830 */
860
866 void
870 void
872 void
874
885 static constexpr std::size_t kAccountCleanupChunk = 4096;
886
908 template <typename OuterMap, typename BeforeErase>
909 void
911 std::uint64_t seq,
912 HashSet<AccountID> const& accounts,
913 OuterMap& outerMap,
914 BeforeErase&& beforeErase);
915
923 void
925 std::uint64_t seq,
926 HashSet<AccountID> const& accounts,
927 SubInfoMapType& subMap);
928
933 void
935
938
940
941 // Independent lock domains so a long cleanup/publish on one does not stall
942 // the others. Hold at most one at a time; if ever more, order: accountLock_,
943 // bookLock_, mptLock_, streamLock_.
944 //
945 // Deferred-destruction rule (non-recursive mutexes): under bookLock_,
946 // mptLock_, or streamLock_, never let the last InfoSub pointer die inside the
947 // lock - ~InfoSub re-acquires it via unsub* -> self-deadlock. Publishers
948 // collect the locked pointers in a container declared before the lock and
949 // destruct after release (see pubServer / pubBookTransaction /
950 // pubMPTTransaction). accountLock_ is exempt: ~InfoSub offloads account
951 // teardown to scheduleAccountCleanup.
956
958
963
965 boost::asio::steady_timer heartbeatTimer_;
966 boost::asio::steady_timer clusterTimer_;
967 boost::asio::steady_timer accountHistoryTxTimer_;
968
970
973
975
977
988
993
995
997
998 // Used as array indices; converting to enum class would require casts at ~40 call sites.
999 // NOLINTNEXTLINE(cppcoreguidelines-use-enum-class)
1001 SLedger, // Accepted ledgers.
1002 SManifests, // Received validator manifests.
1003 SServer, // When server changes connectivity state.
1004 STransactions, // All accepted transactions.
1005 SRtTransactions, // All proposed and accepted transactions.
1006 SValidations, // Received validations.
1007 SPeerStatus, // Peer status changes.
1008 SConsensusPhase, // Consensus phase
1009 SBookChanges, // Per-ledger order book changes
1010 SLastEntry // Any new entry must be ADDED ABOVE this one
1011 };
1012
1018
1020
1022
1023 // Whether we are in standalone mode.
1024 bool const standalone_;
1025
1026 // The number of nodes that we need to consider ourselves connected.
1028
1029 // Transaction batching.
1034
1036
1039
1040private:
1041 struct Stats
1042 {
1043 template <class Handler>
1044 Stats(Handler const& handler, beast::insight::Collector::Ptr const& collector)
1045 : hook(collector->makeHook(handler))
1047 collector->makeGauge("State_Accounting", "Disconnected_duration"))
1048 , connectedDuration(collector->makeGauge("State_Accounting", "Connected_duration"))
1049 , syncingDuration(collector->makeGauge("State_Accounting", "Syncing_duration"))
1050 , trackingDuration(collector->makeGauge("State_Accounting", "Tracking_duration"))
1051 , fullDuration(collector->makeGauge("State_Accounting", "Full_duration"))
1053 collector->makeGauge("State_Accounting", "Disconnected_transitions"))
1055 collector->makeGauge("State_Accounting", "Connected_transitions"))
1056 , syncingTransitions(collector->makeGauge("State_Accounting", "Syncing_transitions"))
1057 , trackingTransitions(collector->makeGauge("State_Accounting", "Tracking_transitions"))
1058 , fullTransitions(collector->makeGauge("State_Accounting", "Full_transitions"))
1059 {
1060 }
1061
1068
1074 };
1075
1076 std::mutex statsMutex_; // Mutex to lock stats_
1078
1079private:
1080 void
1082};
1083
1084//------------------------------------------------------------------------------
1085
1087 {"disconnected", "connected", "syncing", "tracking", "full"}};
1088
1090
1097
1098static auto const kGenesisAccountId =
1099 calcAccountID(generateKeyPair(KeyType::Secp256k1, generateSeed("masterpassphrase")).first);
1100
1101//------------------------------------------------------------------------------
1102inline OperatingMode
1104{
1105 return mode_;
1106}
1107
1108inline std::string
1109NetworkOPsImp::strOperatingMode(bool const admin /* = false */) const
1110{
1111 return strOperatingMode(mode_, admin);
1112}
1113
1114inline void
1119
1120inline void
1125
1126inline void
1131
1132inline bool
1137
1138inline bool
1143
1146{
1147 static std::string const kHostname = boost::asio::ip::host_name();
1148
1149 if (forAdmin)
1150 return kHostname;
1151
1152 // For non-admin uses hash the node public key into a
1153 // single RFC1751 word:
1154 static std::string const kShroudedHostId = [this]() {
1155 auto const& id = registry_.get().getApp().nodeIdentity();
1156
1157 return RFC1751::getWordFromBlob(id.first.data(), id.first.size());
1158 }();
1159
1160 return kShroudedHostId;
1161}
1162
1163void
1165{
1167
1168 // Only do this work if a cluster is configured
1169 if (registry_.get().getCluster().size() != 0)
1171}
1172
1173void
1175 boost::asio::steady_timer& timer,
1176 std::chrono::milliseconds const& expiryTime,
1177 std::function<void()> onExpire,
1178 std::function<void()> onError)
1179{
1180 // Only start the timer if waitHandlerCounter_ is not yet joined.
1181 if (auto optionalCountedHandler =
1182 waitHandlerCounter_.wrap([this, onExpire, onError](boost::system::error_code const& e) {
1183 if ((e.value() == boost::system::errc::success) && (!jobQueue_.isStopped()))
1184 {
1185 onExpire();
1186 }
1187 // Recover as best we can if an unexpected error occurs.
1188 if (e.value() != boost::system::errc::success &&
1189 e.value() != boost::asio::error::operation_aborted)
1190 {
1191 // Try again later and hope for the best.
1192 JLOG(journal_.error())
1193 << "Timer got error '" << e.message() << "'. Restarting timer.";
1194 onError();
1195 }
1196 }))
1197 {
1198 timer.expires_after(expiryTime);
1199 timer.async_wait(std::move(*optionalCountedHandler));
1200 }
1201}
1202
1203void
1205{
1206 setTimer(
1208 consensus_.parms().ledgerGRANULARITY,
1209 [this]() {
1210 jobQueue_.addJob(JtNetopTimer, "NetHeart", [this]() { processHeartbeatTimer(); });
1211 },
1212 [this]() { setHeartbeatTimer(); });
1213}
1214
1215void
1217{
1218 using namespace std::chrono_literals;
1219
1220 setTimer(
1222 10s,
1223 [this]() {
1224 jobQueue_.addJob(JtNetopCluster, "NetCluster", [this]() { processClusterTimer(); });
1225 },
1226 [this]() { setClusterTimer(); });
1227}
1228
1229void
1231{
1232 JLOG(journal_.debug()) << "Scheduling AccountHistory job for account "
1233 << toBase58(subInfo.index->accountId);
1234 using namespace std::chrono_literals;
1235 setTimer(
1237 4s,
1238 [this, subInfo]() { addAccountHistoryJob(subInfo); },
1239 [this, subInfo]() { setAccountHistoryJobTimer(subInfo); });
1240}
1241
1242void
1244{
1245 RclConsensusLogger clog("Heartbeat Timer", consensus_.validating(), journal_);
1246 {
1247 std::unique_lock lock{registry_.get().getApp().getMasterMutex()};
1248
1249 // VFALCO NOTE This is for diagnosing a crash on exit
1250 LoadManager& mgr(registry_.get().getLoadManager());
1251 mgr.heartbeat();
1252
1253 std::size_t const numPeers = registry_.get().getOverlay().size();
1254
1255 // do we have sufficient peers? If not, we are disconnected.
1256 if (numPeers < minPeerCount_)
1257 {
1259 {
1262 ss << "Node count (" << numPeers << ") has fallen "
1263 << "below required minimum (" << minPeerCount_ << ").";
1264 JLOG(journal_.warn()) << ss.str();
1265 CLOG(clog.ss()) << "set mode to DISCONNECTED: " << ss.str();
1266 }
1267 else
1268 {
1269 CLOG(clog.ss()) << "already DISCONNECTED. too few peers (" << numPeers
1270 << "), need at least " << minPeerCount_;
1271 }
1272
1273 // MasterMutex lock need not be held to call setHeartbeatTimer()
1274 lock.unlock();
1275 // We do not call consensus_.timerEntry until there are enough
1276 // peers providing meaningful inputs to consensus
1278
1279 return;
1280 }
1281
1283 {
1285 JLOG(journal_.info()) << "Node count (" << numPeers << ") is sufficient.";
1286 CLOG(clog.ss()) << "setting mode to CONNECTED based on " << numPeers << " peers. ";
1287 }
1288
1289 // Check if the last validated ledger forces a change between these
1290 // states.
1291 auto origMode = mode_.load();
1292 CLOG(clog.ss()) << "mode: " << strOperatingMode(origMode, true);
1294 {
1296 }
1297 else if (mode_ == OperatingMode::CONNECTED)
1298 {
1300 }
1301 auto newMode = mode_.load();
1302 if (origMode != newMode)
1303 {
1304 CLOG(clog.ss()) << ", changing to " << strOperatingMode(newMode, true);
1305 }
1306 CLOG(clog.ss()) << ". ";
1307 }
1308
1309 consensus_.timerEntry(registry_.get().getTimeKeeper().closeTime(), clog.ss());
1310
1311 CLOG(clog.ss()) << "consensus phase " << to_string(lastConsensusPhase_);
1312 ConsensusPhase const currPhase = consensus_.phase();
1313 if (lastConsensusPhase_ != currPhase)
1314 {
1315 reportConsensusStateChange(currPhase);
1316 lastConsensusPhase_ = currPhase;
1317 CLOG(clog.ss()) << " changed to " << to_string(lastConsensusPhase_);
1318 }
1319 CLOG(clog.ss()) << ". ";
1320
1322}
1323
1324void
1326{
1327 if (registry_.get().getCluster().size() == 0)
1328 return;
1329
1330 using namespace std::chrono_literals;
1331
1332 bool const update = registry_.get().getCluster().update(
1333 registry_.get().getApp().nodeIdentity().first,
1334 "",
1335 (ledgerMaster_.getValidatedLedgerAge() <= 4min)
1336 ? registry_.get().getFeeTrack().getLocalFee()
1337 : 0,
1338 registry_.get().getTimeKeeper().now());
1339
1340 if (!update)
1341 {
1342 JLOG(journal_.debug()) << "Too soon to send cluster update";
1344 return;
1345 }
1346
1347 protocol::TMCluster cluster;
1348 registry_.get().getCluster().forEach([&cluster](ClusterNode const& node) {
1349 protocol::TMClusterNode& n = *cluster.add_clusternodes();
1350 n.set_publickey(toBase58(TokenType::NodePublic, node.identity()));
1351 n.set_reporttime(node.getReportTime().time_since_epoch().count());
1352 n.set_nodeload(node.getLoadFee());
1353 if (!node.name().empty())
1354 n.set_nodename(node.name());
1355 });
1356
1357 resource::Gossip const gossip = registry_.get().getResourceManager().exportConsumers();
1358 for (auto& item : gossip.items)
1359 {
1360 protocol::TMLoadSource& node = *cluster.add_loadsources();
1361 node.set_name(to_string(item.address));
1362 node.set_cost(item.balance);
1363 }
1364 registry_.get().getOverlay().foreach(
1365 sendIf(std::make_shared<Message>(cluster, protocol::mtCLUSTER), PeerInCluster()));
1367}
1368
1369//------------------------------------------------------------------------------
1370
1372NetworkOPsImp::strOperatingMode(OperatingMode const mode, bool const admin) const
1373{
1374 if (mode == OperatingMode::FULL && admin)
1375 {
1376 auto const consensusMode = consensus_.mode();
1377 if (consensusMode != ConsensusMode::WrongLedger)
1378 {
1379 if (consensusMode == ConsensusMode::Proposing)
1380 return "proposing";
1381
1382 if (consensus_.validating())
1383 return "validating";
1384 }
1385 }
1386
1387 return kStates[static_cast<std::size_t>(mode)];
1388}
1389
1390void
1392{
1393 if (isNeedNetworkLedger())
1394 {
1395 // Nothing we can do if we've never been in sync
1396 return;
1397 }
1398
1399 // Reject any transaction with the tfInnerBatchTxn flag at the network
1400 // boundary, regardless of amendment state.
1401 if (iTrans->isFlag(tfInnerBatchTxn))
1402 {
1403 JLOG(journal_.error()) << "Submitted transaction invalid: tfInnerBatchTxn flag present.";
1404 return;
1405 }
1406
1407 // this is an asynchronous interface
1408 auto const trans = sterilize(*iTrans);
1409
1410 auto const txid = trans->getTransactionID();
1411 auto const flags = registry_.get().getHashRouter().getFlags(txid);
1412
1414 {
1415 JLOG(journal_.warn()) << "Submitted transaction cached bad";
1416 return;
1417 }
1418
1419 try
1420 {
1421 auto const [validity, reason] = checkValidity(
1422 registry_.get().getHashRouter(), *trans, ledgerMaster_.getValidatedRules());
1423
1424 if (validity != Validity::Valid)
1425 {
1426 JLOG(journal_.warn()) << "Submitted transaction invalid: " << reason;
1427 return;
1428 }
1429 }
1430 catch (std::exception const& ex)
1431 {
1432 JLOG(journal_.warn()) << "Exception checking transaction " << txid << ": " << ex.what();
1433
1434 return;
1435 }
1436
1437 std::string reason;
1438
1439 auto tx = std::make_shared<Transaction>(trans, reason, registry_.get().getApp());
1440
1441 jobQueue_.addJob(JtTransaction, "SubmitTxn", [this, tx]() {
1442 auto t = tx;
1443 processTransaction(t, false, false, FailHard::No);
1444 });
1445}
1446
1447bool
1449{
1450 auto const newFlags = registry_.get().getHashRouter().getFlags(transaction->getID());
1451
1453 {
1454 // cached bad
1455 JLOG(journal_.warn()) << transaction->getID() << ": cached bad!\n";
1456 transaction->setStatus(TransStatus::INVALID);
1457 transaction->setResult(temBAD_SIGNATURE);
1458 return false;
1459 }
1460
1461 auto const view = ledgerMaster_.getCurrentLedger();
1462
1463 // This function is called by several different parts of the codebase
1464 // under no circumstances will we ever accept an inner txn within a batch
1465 // txn from the network.
1466 auto const sttx = *transaction->getSTransaction();
1467 if (sttx.isFlag(tfInnerBatchTxn))
1468 {
1469 transaction->setStatus(TransStatus::INVALID);
1470 transaction->setResult(temINVALID_FLAG);
1471 registry_.get().getHashRouter().setFlags(transaction->getID(), HashRouterFlags::BAD);
1472 return false;
1473 }
1474
1475 // NOTE ximinez - I think this check is redundant,
1476 // but I'm not 100% sure yet.
1477 //
1478 // For an ordinary transaction it is: the relay and submit paths have
1479 // already run checkValidity, so this costs a HashRouter lookup. It is not
1480 // redundant for a role-signature transaction while fixCleanup3_4_0 is
1481 // activating. Those paths verify against the validated rules, which lag
1482 // the open ledger rules used here, and checkValidity scopes a cached
1483 // verdict to the rules that reached it, so this call can verify the
1484 // signature again and come to a different answer. SigBad is therefore
1485 // reachable, and the handler below is the correct response to it.
1486 auto const& viewRules = view->rules();
1487 auto const [validity, reason] = checkValidity(registry_.get().getHashRouter(), sttx, viewRules);
1488
1489 // Not concerned with local checks at this point.
1490 if (validity == Validity::SigBad)
1491 {
1492 JLOG(journal_.info()) << "Transaction has bad signature: " << reason;
1493 transaction->setStatus(TransStatus::INVALID);
1494 transaction->setResult(temBAD_SIGNATURE);
1495 // See the matching guard in PeerImp::checkTransaction: only cache
1496 // BAD for a role-signature transaction once fixCleanup3_4_0 is
1497 // enabled on this node. Remove together with the amendment.
1498 if (viewRules.enabled(fixCleanup3_4_0) ||
1499 (!sttx.isFieldPresent(sfSponsorSignature) &&
1500 !sttx.isFieldPresent(sfCounterpartySignature)))
1501 {
1502 registry_.get().getHashRouter().setFlags(transaction->getID(), HashRouterFlags::BAD);
1503 }
1504 return false;
1505 }
1506
1507 // canonicalize can change our pointer
1508 registry_.get().getMasterTransaction().canonicalize(&transaction);
1509
1510 return true;
1511}
1512
1513void
1515 std::shared_ptr<Transaction>& transaction,
1516 bool bUnlimited,
1517 bool bLocal,
1518 FailHard failType)
1519{
1520 auto ev = jobQueue_.makeLoadEvent(JtTxnProc, "ProcessTXN");
1521
1522 // preProcessTransaction can change our pointer
1523 if (!preProcessTransaction(transaction))
1524 return;
1525
1526 if (bLocal)
1527 {
1528 doTransactionSync(transaction, bUnlimited, failType);
1529 }
1530 else
1531 {
1532 doTransactionAsync(transaction, bUnlimited, failType);
1533 }
1534}
1535
1536void
1538 std::shared_ptr<Transaction> transaction,
1539 bool bUnlimited,
1540 FailHard failType)
1541{
1542 std::scoped_lock const lock(mutex_);
1543
1544 if (transaction->getApplying())
1545 return;
1546
1547 transactions_.emplace_back(transaction, bUnlimited, false, failType);
1548 transaction->setApplying();
1549
1551 {
1552 if (jobQueue_.addJob(JtBatch, "TxBatchAsync", [this]() { transactionBatch(); }))
1553 {
1555 }
1556 }
1557}
1558
1559void
1561 std::shared_ptr<Transaction> transaction,
1562 bool bUnlimited,
1563 FailHard failType)
1564{
1566
1567 if (!transaction->getApplying())
1568 {
1569 transactions_.emplace_back(transaction, bUnlimited, true, failType);
1570 transaction->setApplying();
1571 }
1572
1573 doTransactionSyncBatch(lock, [&transaction](std::unique_lock<std::mutex> const&) {
1574 return transaction->getApplying();
1575 });
1576}
1577
1578void
1581 std::function<bool(std::unique_lock<std::mutex> const&)> retryCallback)
1582{
1583 do
1584 {
1586 {
1587 // A batch processing job is already running, so wait.
1588 cond_.wait(lock);
1589 }
1590 else
1591 {
1592 apply(lock);
1593
1594 if (!transactions_.empty())
1595 {
1596 // More transactions need to be applied, but by another job.
1597 if (jobQueue_.addJob(JtBatch, "TxBatchSync", [this]() { transactionBatch(); }))
1598 {
1600 }
1601 }
1602 }
1603 } while (retryCallback(lock));
1604}
1605
1606void
1608{
1609 auto ev = jobQueue_.makeLoadEvent(JtTxnProc, "ProcessTXNSet");
1611 candidates.reserve(set.size());
1612 for (auto const& [_, tx] : set)
1613 {
1614 std::string reason;
1615 auto transaction = std::make_shared<Transaction>(tx, reason, registry_.get().getApp());
1616
1617 if (transaction->getStatus() == TransStatus::INVALID)
1618 {
1619 if (!reason.empty())
1620 {
1621 JLOG(journal_.trace()) << "Exception checking transaction: " << reason;
1622 }
1623 registry_.get().getHashRouter().setFlags(tx->getTransactionID(), HashRouterFlags::BAD);
1624 continue;
1625 }
1626
1627 // preProcessTransaction can change our pointer
1628 if (!preProcessTransaction(transaction))
1629 continue;
1630
1631 candidates.emplace_back(transaction);
1632 }
1633
1635 transactions.reserve(candidates.size());
1636
1638
1639 for (auto& transaction : candidates)
1640 {
1641 if (!transaction->getApplying())
1642 {
1643 transactions.emplace_back(transaction, false, false, FailHard::No);
1644 transaction->setApplying();
1645 }
1646 }
1647
1648 if (transactions_.empty())
1649 {
1651 }
1652 else
1653 {
1654 transactions_.reserve(transactions_.size() + transactions.size());
1655 for (auto& t : transactions)
1656 transactions_.push_back(std::move(t));
1657 }
1658 if (transactions_.empty())
1659 {
1660 JLOG(journal_.debug()) << "No transaction to process!";
1661 return;
1662 }
1663
1665 XRPL_ASSERT(lock.owns_lock(), "xrpl::NetworkOPsImp::processTransactionSet has lock");
1666 return std::ranges::any_of(
1667 transactions_, [](auto const& t) { return t.transaction->getApplying(); });
1668 });
1669}
1670
1671void
1673{
1675
1677 return;
1678
1679 while (!transactions_.empty())
1680 {
1681 apply(lock);
1682 }
1683}
1684
1685void
1687{
1691 XRPL_ASSERT(!transactions.empty(), "xrpl::NetworkOPsImp::apply : non-empty transactions");
1692 XRPL_ASSERT(
1693 dispatchState_ != DispatchState::Running, "xrpl::NetworkOPsImp::apply : is not running");
1694
1696
1697 batchLock.unlock();
1698
1699 {
1700 std::unique_lock masterLock{registry_.get().getApp().getMasterMutex(), std::defer_lock};
1701 bool changed = false;
1702 {
1703 std::unique_lock ledgerLock{ledgerMaster_.peekMutex(), std::defer_lock};
1704 std::lock(masterLock, ledgerLock);
1705 registry_.get().getOpenLedger().modify([&](OpenView& view, beast::Journal j) {
1707 {
1708 // we check before adding to the batch
1709 ApplyFlags flags = TapNone;
1710 if (e.admin)
1711 flags |= TapUnlimited;
1712
1713 if (e.failType == FailHard::Yes)
1714 flags |= TapFailHard;
1715
1716 auto const result = registry_.get().getTxQ().apply(
1717 registry_.get().getApp(), view, e.transaction->getSTransaction(), flags, j);
1718 e.result = result.ter;
1719 e.applied = result.applied;
1720 changed = changed || result.applied;
1721 }
1722 return changed;
1723 });
1724 }
1725 if (changed)
1727
1728 std::optional<LedgerIndex> validatedLedgerIndex;
1729 if (auto const l = ledgerMaster_.getValidatedLedger())
1730 validatedLedgerIndex = l->header().seq;
1731
1732 auto newOL = registry_.get().getOpenLedger().current();
1733 for (TransactionStatus const& e : transactions)
1734 {
1735 e.transaction->clearSubmitResult();
1736
1737 if (e.applied)
1738 {
1739 pubProposedTransaction(newOL, e.transaction->getSTransaction(), e.result);
1740 e.transaction->setApplied();
1741 }
1742
1743 e.transaction->setResult(e.result);
1744
1745 if (isTemMalformed(e.result))
1746 {
1747 registry_.get().getHashRouter().setFlags(
1748 e.transaction->getID(), HashRouterFlags::BAD);
1749 }
1750
1751#ifdef DEBUG
1752 if (!isTesSuccess(e.result))
1753 {
1754 std::string token, human;
1755
1756 if (transResultInfo(e.result, token, human))
1757 {
1758 JLOG(journal_.info()) << "TransactionResult: " << token << ": " << human;
1759 }
1760 }
1761#endif
1762
1763 bool const addLocal = e.local;
1764
1765 if (isTesSuccess(e.result))
1766 {
1767 JLOG(journal_.debug()) << "Transaction is now included in open ledger";
1768 e.transaction->setStatus(TransStatus::INCLUDED);
1769
1770 // Pop as many "reasonable" transactions for this account as
1771 // possible. "Reasonable" means they have sequential sequence
1772 // numbers, or use tickets.
1773 auto const& txCur = e.transaction->getSTransaction();
1774
1775 std::size_t count = 0;
1776 for (auto txNext = ledgerMaster_.popAcctTransaction(txCur);
1777 txNext && count < kMaxPoppedTransactions;
1778 txNext = ledgerMaster_.popAcctTransaction(txCur), ++count)
1779 {
1780 if (!batchLock.owns_lock())
1781 batchLock.lock();
1782 std::string reason;
1783 auto const trans = sterilize(*txNext);
1784 auto t = std::make_shared<Transaction>(trans, reason, registry_.get().getApp());
1785 if (t->getApplying())
1786 break;
1787 submitHeld.emplace_back(t, false, false, FailHard::No);
1788 t->setApplying();
1789 }
1790 if (batchLock.owns_lock())
1791 batchLock.unlock();
1792 }
1793 else if (e.result == tefPAST_SEQ)
1794 {
1795 // duplicate or conflict
1796 JLOG(journal_.info()) << "Transaction is obsolete";
1797 e.transaction->setStatus(TransStatus::OBSOLETE);
1798 }
1799 else if (e.result == terQUEUED)
1800 {
1801 JLOG(journal_.debug()) << "Transaction is likely to claim a"
1802 << " fee, but is queued until fee drops";
1803
1804 e.transaction->setStatus(TransStatus::HELD);
1805 // Add to held transactions, because it could get
1806 // kicked out of the queue, and this will try to
1807 // put it back.
1808 ledgerMaster_.addHeldTransaction(e.transaction);
1809 e.transaction->setQueued();
1810 e.transaction->setKept();
1811 }
1812 else if (isTerRetry(e.result) || isTelLocal(e.result) || isTefFailure(e.result))
1813 {
1814 if (e.failType != FailHard::Yes)
1815 {
1816 auto const lastLedgerSeq =
1817 e.transaction->getSTransaction()->at(~sfLastLedgerSequence);
1818 auto const ledgersLeft = lastLedgerSeq
1819 ? *lastLedgerSeq - ledgerMaster_.getCurrentLedgerIndex()
1821 // If any of these conditions are met, the transaction can
1822 // be held:
1823 // 1. It was submitted locally. (Note that this flag is only
1824 // true on the initial submission.)
1825 // 2. The transaction has a LastLedgerSequence, and the
1826 // LastLedgerSequence is fewer than LocalTxs::kHoldLedgers
1827 // (5) ledgers into the future. (Remember that an
1828 // unseated optional compares as less than all seated
1829 // values, so it has to be checked explicitly first.)
1830 // 3. The HashRouterFlags::BAD flag is not set on the txID.
1831 // (setFlags
1832 // checks before setting. If the flag is set, it returns
1833 // false, which means it's been held once without one of
1834 // the other conditions, so don't hold it again. Time's
1835 // up!)
1836 //
1837 if (e.local || (ledgersLeft && ledgersLeft <= LocalTxs::kHoldLedgers) ||
1838 registry_.get().getHashRouter().setFlags(
1839 e.transaction->getID(), HashRouterFlags::HELD))
1840 {
1841 // transaction should be held
1842 JLOG(journal_.debug()) << "Transaction should be held: " << e.result;
1843 e.transaction->setStatus(TransStatus::HELD);
1844 ledgerMaster_.addHeldTransaction(e.transaction);
1845 e.transaction->setKept();
1846 }
1847 else
1848 {
1849 JLOG(journal_.debug())
1850 << "Not holding transaction " << e.transaction->getID() << ": "
1851 << (e.local ? "local" : "network") << ", "
1852 << "result: " << e.result << " ledgers left: "
1853 << (ledgersLeft ? to_string(*ledgersLeft) : "unspecified");
1854 }
1855 }
1856 }
1857 else
1858 {
1859 JLOG(journal_.debug()) << "Status other than success " << e.result;
1860 e.transaction->setStatus(TransStatus::INVALID);
1861 }
1862
1863 auto const enforceFailHard = e.failType == FailHard::Yes && !isTesSuccess(e.result);
1864
1865 if (addLocal && !enforceFailHard)
1866 {
1867 localTX_->pushBack(
1868 ledgerMaster_.getCurrentLedgerIndex(), e.transaction->getSTransaction());
1869 e.transaction->setKept();
1870 }
1871
1872 if ((e.applied ||
1873 ((mode_ != OperatingMode::FULL) && (e.failType != FailHard::Yes) && e.local) ||
1874 (e.result == terQUEUED)) &&
1875 !enforceFailHard)
1876 {
1877 auto const toSkip =
1878 registry_.get().getHashRouter().shouldRelay(e.transaction->getID());
1879 if (auto const sttx = *(e.transaction->getSTransaction()); toSkip &&
1880 // Skip relaying if it's an inner batch txn. The flag should
1881 // only be set if the Batch feature is enabled. If Batch is
1882 // not enabled, the flag is always invalid, so don't relay
1883 // it regardless.
1884 !(sttx.isFlag(tfInnerBatchTxn)))
1885 {
1886 protocol::TMTransaction tx;
1887 Serializer s;
1888
1889 sttx.add(s);
1890 tx.set_rawtransaction(s.data(), s.size());
1891 tx.set_status(protocol::tsCURRENT);
1892 tx.set_receivetimestamp(
1893 registry_.get().getTimeKeeper().now().time_since_epoch().count());
1894 tx.set_deferred(e.result == terQUEUED);
1895 // FIXME: This should be when we received it
1896 registry_.get().getOverlay().relay(e.transaction->getID(), tx, *toSkip);
1897 e.transaction->setBroadcast();
1898 }
1899 }
1900
1901 if (validatedLedgerIndex)
1902 {
1903 auto maybeFeeAndSeq = registry_.get().getTxQ().getTxRequiredFeeAndSeq(
1904 *newOL, e.transaction->getSTransaction());
1905 if (maybeFeeAndSeq.has_value())
1906 {
1907 auto [fee, accountSeq, availableSeq] = *maybeFeeAndSeq;
1908 e.transaction->setCurrentLedgerState(
1909 *validatedLedgerIndex, fee, accountSeq, availableSeq);
1910 }
1911 else
1912 {
1913 JLOG(journal_.debug())
1914 << "Unable to compute current ledger state for tx "
1915 << e.transaction->getID() << " in validated ledger "
1916 << *validatedLedgerIndex << ": " << transToken(maybeFeeAndSeq.error());
1917 }
1918 }
1919 }
1920 }
1921
1922 batchLock.lock();
1923
1924 for (TransactionStatus const& e : transactions)
1925 e.transaction->clearApplying();
1926
1927 if (!submitHeld.empty())
1928 {
1929 if (transactions_.empty())
1930 {
1931 transactions_.swap(submitHeld);
1932 }
1933 else
1934 {
1935 transactions_.reserve(transactions_.size() + submitHeld.size());
1936 for (auto& e : submitHeld)
1937 transactions_.push_back(std::move(e));
1938 }
1939 }
1940
1941 cond_.notify_all();
1942
1944}
1945
1946//
1947// Owner functions
1948//
1949
1952{
1954 auto root = keylet::ownerDir(account);
1955 auto sleNode = lpLedger->read(keylet::page(root));
1956 if (sleNode)
1957 {
1958 std::uint64_t uNodeDir = 0;
1959
1960 do
1961 {
1962 for (auto const& uDirEntry : sleNode->getFieldV256(sfIndexes))
1963 {
1964 auto sleCur = lpLedger->read(keylet::child(uDirEntry));
1965 XRPL_ASSERT(sleCur, "xrpl::NetworkOPsImp::getOwnerInfo : non-null child SLE");
1966
1967 switch (sleCur->getType())
1968 {
1969 case ltOFFER:
1970 if (!jvObjects.isMember(jss::offers))
1971 jvObjects[jss::offers] = json::Value(json::ValueType::Array);
1972
1973 jvObjects[jss::offers].append(sleCur->getJson(JsonOptions::Values::None));
1974 break;
1975
1976 case ltRIPPLE_STATE:
1977 if (!jvObjects.isMember(jss::ripple_lines))
1978 {
1979 jvObjects[jss::ripple_lines] = json::Value(json::ValueType::Array);
1980 }
1981
1982 jvObjects[jss::ripple_lines].append(
1983 sleCur->getJson(JsonOptions::Values::None));
1984 break;
1985
1986 case ltACCOUNT_ROOT:
1987 case ltDIR_NODE:
1988 // LCOV_EXCL_START
1989 default:
1990 UNREACHABLE(
1991 "xrpl::NetworkOPsImp::getOwnerInfo : invalid "
1992 "type");
1993 break;
1994 // LCOV_EXCL_STOP
1995 }
1996 }
1997
1998 uNodeDir = sleNode->getFieldU64(sfIndexNext);
1999
2000 if (uNodeDir != 0u)
2001 {
2002 sleNode = lpLedger->read(keylet::page(root, uNodeDir));
2003 XRPL_ASSERT(sleNode, "xrpl::NetworkOPsImp::getOwnerInfo : read next page");
2004 }
2005 } while (uNodeDir != 0u);
2006 }
2007
2008 return jvObjects;
2009}
2010
2011//
2012// Other
2013//
2014
2015inline bool
2020
2021inline bool
2026
2027void
2033
2034inline bool
2039
2040inline void
2045
2046inline void
2051
2052inline bool
2054{
2055 return unlBlocked_;
2056}
2057
2058void
2064
2065inline void
2070
2071bool
2073{
2074 // Returns true if there's an *abnormal* ledger issue, normal changing in
2075 // TRACKING mode should return false. Do we have sufficient validations for
2076 // our last closed ledger? Or do sufficient nodes agree? And do we have no
2077 // better ledger available? If so, we are either tracking or full.
2078
2079 JLOG(journal_.trace()) << "NetworkOPsImp::checkLastClosedLedger";
2080
2081 auto const ourClosed = ledgerMaster_.getClosedLedger();
2082
2083 if (!ourClosed)
2084 return false;
2085
2086 UInt256 closedLedger = ourClosed->header().hash;
2087 UInt256 const prevClosedLedger = ourClosed->header().parentHash;
2088 JLOG(journal_.trace()) << "OurClosed: " << closedLedger;
2089 JLOG(journal_.trace()) << "PrevClosed: " << prevClosedLedger;
2090
2091 //-------------------------------------------------------------------------
2092 // Determine preferred last closed ledger
2093
2094 auto& validations = registry_.get().getValidations();
2095 JLOG(journal_.debug()) << "ValidationTrie " << json::Compact(validations.getJsonTrie());
2096
2097 // Will rely on peer LCL if no trusted validations exist
2099 peerCounts[closedLedger] = 0;
2101 peerCounts[closedLedger]++;
2102
2103 for (auto& peer : peerList)
2104 {
2105 UInt256 const peerLedger = peer->getClosedLedgerHash();
2106
2107 if (peerLedger.isNonZero())
2108 ++peerCounts[peerLedger];
2109 }
2110
2111 for (auto const& it : peerCounts)
2112 JLOG(journal_.debug()) << "L: " << it.first << " n=" << it.second;
2113
2114 UInt256 const preferredLCL = validations.getPreferredLCL(
2115 RCLValidatedLedger{ourClosed, validations.adaptor().journal()},
2116 ledgerMaster_.getValidLedgerIndex(),
2117 peerCounts);
2118
2119 bool switchLedgers = preferredLCL != closedLedger;
2120 if (switchLedgers)
2121 closedLedger = preferredLCL;
2122 //-------------------------------------------------------------------------
2123 if (switchLedgers && (closedLedger == prevClosedLedger))
2124 {
2125 // don't switch to our own previous ledger
2126 JLOG(journal_.info()) << "We won't switch to our own previous ledger";
2127 networkClosed = ourClosed->header().hash;
2128 switchLedgers = false;
2129 }
2130 else
2131 {
2132 networkClosed = closedLedger;
2133 }
2134
2135 if (!switchLedgers)
2136 return false;
2137
2138 auto consensus = ledgerMaster_.getLedgerByHash(closedLedger);
2139
2140 if (!consensus)
2141 {
2142 consensus = registry_.get().getInboundLedgers().acquire(
2143 closedLedger, 0, InboundLedger::Reason::CONSENSUS);
2144 }
2145
2146 if (consensus &&
2147 (!ledgerMaster_.canBeCurrent(consensus) ||
2148 !ledgerMaster_.isCompatible(*consensus, journal_.debug(), "Not switching")))
2149 {
2150 // Don't switch to a ledger not on the validated chain
2151 // or with an invalid close time or sequence
2152 networkClosed = ourClosed->header().hash;
2153 return false;
2154 }
2155
2156 JLOG(journal_.warn()) << "We are not running on the consensus ledger";
2157 JLOG(journal_.info()) << "Our LCL: " << ourClosed->header().hash << getJson({*ourClosed, {}});
2158 JLOG(journal_.info()) << "Net LCL " << closedLedger;
2159
2161 {
2163 }
2164
2165 if (consensus)
2166 {
2167 // FIXME: If this rewinds the ledger sequence, or has the same
2168 // sequence, we should update the status on any stored transactions
2169 // in the invalidated ledgers.
2170 switchLastClosedLedger(consensus);
2171 }
2172
2173 return true;
2174}
2175
2176void
2178{
2179 // set the newLCL as our last closed ledger -- this is abnormal code
2180 JLOG(journal_.error()) << "JUMP last closed ledger to " << newLCL->header().hash;
2181
2183
2184 // Update fee computations.
2185 registry_.get().getTxQ().processClosedLedger(registry_.get().getApp(), *newLCL, true);
2186
2187 // Caller must own master lock
2188 {
2189 // Apply tx in old open ledger to new
2190 // open ledger. Then apply local tx.
2191
2192 auto retries = localTX_->getTxSet();
2193 auto const lastVal = registry_.get().getLedgerMaster().getValidatedLedger();
2195 if (lastVal)
2196 {
2197 rules = makeRulesGivenLedger(*lastVal, registry_.get().getApp().config().features);
2198 }
2199 else
2200 {
2201 rules.emplace(registry_.get().getApp().config().features);
2202 }
2203 registry_.get().getOpenLedger().accept(
2204 registry_.get().getApp(),
2205 *rules,
2206 newLCL,
2207 OrderedTxs({}),
2208 false,
2209 retries,
2210 TapNone,
2211 "jump",
2212 [&](OpenView& view, beast::Journal j) {
2213 // Stuff the ledger with transactions from the queue.
2214 return registry_.get().getTxQ().accept(registry_.get().getApp(), view);
2215 });
2216 }
2217
2218 ledgerMaster_.switchLCL(newLCL);
2219
2220 protocol::TMStatusChange s;
2221 s.set_newevent(protocol::neSWITCHED_LEDGER);
2222 s.set_ledgerseq(newLCL->header().seq);
2223 s.set_networktime(registry_.get().getTimeKeeper().now().time_since_epoch().count());
2224 s.set_ledgerhashprevious(
2225 newLCL->header().parentHash.begin(), newLCL->header().parentHash.size());
2226 s.set_ledgerhash(newLCL->header().hash.begin(), newLCL->header().hash.size());
2227 registry_.get().getOverlay().foreach(
2228 SendAlways(std::make_shared<Message>(s, protocol::mtSTATUS_CHANGE)));
2229}
2230
2231bool
2233 UInt256 const& networkClosed,
2235{
2236 XRPL_ASSERT(networkClosed.isNonZero(), "xrpl::NetworkOPsImp::beginConsensus : nonzero input");
2237
2238 auto closingInfo = ledgerMaster_.getCurrentLedger()->header();
2239
2240 JLOG(journal_.info()) << "Consensus time for #" << closingInfo.seq << " with LCL "
2241 << closingInfo.parentHash;
2242
2243 auto prevLedger = ledgerMaster_.getLedgerByHash(closingInfo.parentHash);
2244
2245 if (!prevLedger)
2246 {
2247 // this shouldn't happen unless we jump ledgers
2249 {
2250 JLOG(journal_.warn()) << "Don't have LCL, going to tracking";
2252 CLOG(clog) << "beginConsensus Don't have LCL, going to tracking. ";
2253 }
2254
2255 CLOG(clog) << "beginConsensus no previous ledger. ";
2256 return false;
2257 }
2258
2259 XRPL_ASSERT(
2260 prevLedger->header().hash == closingInfo.parentHash,
2261 "xrpl::NetworkOPsImp::beginConsensus : prevLedger hash matches "
2262 "parent");
2263 XRPL_ASSERT(
2264 closingInfo.parentHash == ledgerMaster_.getClosedLedger()->header().hash,
2265 "xrpl::NetworkOPsImp::beginConsensus : closedLedger parent matches "
2266 "hash");
2267
2268 registry_.get().getValidators().setNegativeUNL(prevLedger->negativeUNL());
2269 TrustChanges const changes = registry_.get().getValidators().updateTrusted(
2270 registry_.get().getValidations().getCurrentNodeIDs(),
2271 closingInfo.parentCloseTime,
2272 *this,
2273 registry_.get().getOverlay(),
2274 registry_.get().getHashRouter());
2275
2276 if (!changes.added.empty() || !changes.removed.empty())
2277 {
2278 registry_.get().getValidations().trustChanged(changes.added, changes.removed);
2279 // Update the AmendmentTable so it tracks the current validators.
2280 registry_.get().getAmendmentTable().trustChanged(
2281 registry_.get().getValidators().getQuorumKeys().second);
2282 }
2283
2284 consensus_.startRound(
2285 registry_.get().getTimeKeeper().closeTime(),
2286 networkClosed,
2287 prevLedger,
2288 changes.removed,
2289 changes.added,
2290 clog);
2291
2292 ConsensusPhase const currPhase = consensus_.phase();
2293 if (lastConsensusPhase_ != currPhase)
2294 {
2295 reportConsensusStateChange(currPhase);
2296 lastConsensusPhase_ = currPhase;
2297 }
2298
2299 JLOG(journal_.debug()) << "Initiating consensus engine";
2300 return true;
2301}
2302
2303bool
2305{
2306 auto const& peerKey = peerPos.publicKey();
2307 if (validatorPK_ == peerKey || validatorMasterPK_ == peerKey)
2308 {
2309 // Could indicate a operator misconfiguration where two nodes are
2310 // running with the same validator key configured, so this isn't fatal,
2311 // and it doesn't necessarily indicate peer misbehavior. But since this
2312 // is a trusted message, it could be a very big deal. Either way, we
2313 // don't want to relay the proposal. Note that the byzantine behavior
2314 // detection in handleNewValidation will notify other peers.
2315 //
2316 // Another, innocuous explanation is unusual message routing and delays,
2317 // causing this node to receive its own messages back.
2318 JLOG(journal_.error()) << "Received a proposal signed by MY KEY from a peer. This may "
2319 "indicate a misconfiguration where another node has the same "
2320 "validator key, or may be caused by unusual message routing and "
2321 "delays.";
2322 return false;
2323 }
2324
2325 return consensus_.peerProposal(registry_.get().getTimeKeeper().closeTime(), peerPos);
2326}
2327
2328void
2330{
2331 // We now have an additional transaction set
2332 // Inform peers we have this set
2333 protocol::TMHaveTransactionSet msg;
2334 msg.set_hash(map->getHash().asUInt256().begin(), 256 / 8);
2335 msg.set_status(protocol::tsHAVE);
2336 registry_.get().getOverlay().foreach(
2337 SendAlways(std::make_shared<Message>(msg, protocol::mtHAVE_SET)));
2338
2339 // We acquired it because consensus asked us to
2340 if (fromAcquire)
2341 consensus_.gotTxSet(registry_.get().getTimeKeeper().closeTime(), RCLTxSet{map});
2342}
2343
2344void
2346{
2347 UInt256 const deadLedger = ledgerMaster_.getClosedLedger()->header().parentHash;
2348 for (auto const& it : registry_.get().getOverlay().getActivePeers())
2349 {
2350 if (it && (it->getClosedLedgerHash() == deadLedger))
2351 {
2352 JLOG(journal_.trace()) << "Killing obsolete peer status";
2353 it->cycleStatus();
2354 }
2355 }
2356
2357 UInt256 networkClosed;
2358 bool const ledgerChange =
2359 checkLastClosedLedger(registry_.get().getOverlay().getActivePeers(), networkClosed);
2360
2361 if (networkClosed.isZero())
2362 {
2363 CLOG(clog) << "endConsensus last closed ledger is zero. ";
2364 return;
2365 }
2366
2367 // WRITEME: Unless we are in FULL and in the process of doing a consensus,
2368 // we must count how many nodes share our LCL, how many nodes disagree with
2369 // our LCL, and how many validations our LCL has. We also want to check
2370 // timing to make sure there shouldn't be a newer LCL. We need this
2371 // information to do the next three tests.
2372
2373 if (((mode_ == OperatingMode::CONNECTED) || (mode_ == OperatingMode::SYNCING)) && !ledgerChange)
2374 {
2375 // Count number of peers that agree with us and UNL nodes whose
2376 // validations we have for LCL. If the ledger is good enough, go to
2377 // TRACKING - TODO
2378 if (!needNetworkLedger_)
2380 }
2381
2383 !ledgerChange)
2384 {
2385 // check if the ledger is good enough to go to FULL
2386 // Note: Do not go to FULL if we don't have the previous ledger
2387 // check if the ledger is bad enough to go to CONNECTED -- TODO
2388 auto current = ledgerMaster_.getCurrentLedger();
2389 if (registry_.get().getTimeKeeper().now() <
2390 (current->header().parentCloseTime + 2 * current->header().closeTimeResolution))
2391 {
2393 }
2394 }
2395
2396 beginConsensus(networkClosed, clog);
2397}
2398
2399void
2407
2408void
2410{
2411 // Hold each locked subscriber alive until after streamLock_ is released:
2412 // if this is the last reference, ~InfoSub re-acquires streamLock_ (via its
2413 // unsub* calls), which would self-deadlock on this non-recursive mutex.
2414 // Declared before the lock so it is destroyed after the lock is dropped.
2416
2417 // VFALCO consider std::shared_mutex
2418 std::scoped_lock const sl(streamLock_);
2419
2420 if (!streamMaps_[SManifests].empty())
2421 {
2423
2424 jvObj[jss::type] = "manifestReceived";
2425 jvObj[jss::master_key] = toBase58(TokenType::NodePublic, mo.masterKey);
2426 if (mo.signingKey)
2427 jvObj[jss::signing_key] = toBase58(TokenType::NodePublic, *mo.signingKey);
2428 jvObj[jss::seq] = json::UInt(mo.sequence);
2429 if (auto sig = mo.getSignature())
2430 jvObj[jss::signature] = strHex(*sig);
2431 jvObj[jss::master_signature] = strHex(mo.getMasterSignature());
2432 if (!mo.domain.empty())
2433 jvObj[jss::domain] = mo.domain;
2434 jvObj[jss::manifest] = strHex(mo.serialized);
2435
2436 for (auto i = streamMaps_[SManifests].begin(); i != streamMaps_[SManifests].end();)
2437 {
2438 if (auto p = i->second.lock())
2439 {
2440 p->send(jvObj, true);
2441 toRelease.push_back(std::move(p));
2442 ++i;
2443 }
2444 else
2445 {
2446 i = streamMaps_[SManifests].erase(i);
2447 }
2448 }
2449 }
2450}
2451
2453 XRPAmount fee,
2454 TxQ::Metrics escalationMetrics, // trivially copyable
2455 LoadFeeTrack const& loadFeeTrack)
2456 : loadFactorServer{loadFeeTrack.getLoadFactor()}
2457 , loadBaseServer{loadFeeTrack.getLoadBase()}
2458 , baseFee{fee}
2459 , em{escalationMetrics}
2460{
2461}
2462
2463bool
2465{
2467 baseFee != b.baseFee || em.has_value() != b.em.has_value())
2468 return true;
2469
2470 if (em && b.em)
2471 {
2472 return (
2473 em->minProcessingFeeLevel != b.em->minProcessingFeeLevel ||
2474 em->openLedgerFeeLevel != b.em->openLedgerFeeLevel ||
2475 em->referenceFeeLevel != b.em->referenceFeeLevel);
2476 }
2477
2478 return false;
2479}
2480
2481// Need to cap to uint64 to uint32 due to JSON limitations
2482static std::uint32_t
2484{
2486
2487 return std::min(kMax32, v);
2488};
2489
2490void
2492{
2493 // Hold each locked subscriber alive until after streamLock_ is released; a
2494 // last-reference ~InfoSub would otherwise re-acquire this non-recursive
2495 // mutex and self-deadlock. Declared before the lock, destroyed after it.
2497
2498 // VFALCO TODO Don't hold the lock across calls to send...make a copy of the
2499 // list into a local array while holding the lock then release
2500 // the lock and call send on everyone.
2501 //
2502 std::scoped_lock const sl(streamLock_);
2503
2504 if (!streamMaps_[SServer].empty())
2505 {
2507
2509 registry_.get().getOpenLedger().current()->fees().base,
2510 registry_.get().getTxQ().getMetrics(*registry_.get().getOpenLedger().current()),
2511 registry_.get().getFeeTrack()};
2512
2513 jvObj[jss::type] = "serverStatus";
2514 jvObj[jss::server_status] = strOperatingMode();
2515 jvObj[jss::load_base] = f.loadBaseServer;
2516 jvObj[jss::load_factor_server] = f.loadFactorServer;
2517 jvObj[jss::base_fee] = f.baseFee.jsonClipped();
2518
2519 if (f.em)
2520 {
2521 auto const loadFactor = std::max(
2523 mulDiv(f.em->openLedgerFeeLevel, f.loadBaseServer, f.em->referenceFeeLevel)
2524 .value_or(xrpl::kMuldivMax));
2525
2526 jvObj[jss::load_factor] = trunc32(loadFactor);
2527 jvObj[jss::load_factor_fee_escalation] = f.em->openLedgerFeeLevel.jsonClipped();
2528 jvObj[jss::load_factor_fee_queue] = f.em->minProcessingFeeLevel.jsonClipped();
2529 jvObj[jss::load_factor_fee_reference] = f.em->referenceFeeLevel.jsonClipped();
2530 }
2531 else
2532 {
2533 jvObj[jss::load_factor] = f.loadFactorServer;
2534 }
2535
2536 lastFeeSummary_ = f;
2537
2538 for (auto i = streamMaps_[SServer].begin(); i != streamMaps_[SServer].end();)
2539 {
2540 InfoSub::pointer p = i->second.lock();
2541
2542 // VFALCO TODO research the possibility of using thread queues and
2543 // linearizing the deletion of subscribers with the
2544 // sending of JSON data.
2545 if (p)
2546 {
2547 p->send(jvObj, true);
2548 toRelease.push_back(std::move(p));
2549 ++i;
2550 }
2551 else
2552 {
2553 i = streamMaps_[SServer].erase(i);
2554 }
2555 }
2556 }
2557}
2558
2559void
2561{
2562 // Hold each locked subscriber alive until after streamLock_ is released; a
2563 // last-reference ~InfoSub would otherwise re-acquire this non-recursive
2564 // mutex and self-deadlock. Declared before the lock, destroyed after it.
2566
2567 std::scoped_lock const sl(streamLock_);
2568
2569 auto& streamMap = streamMaps_[SConsensusPhase];
2570 if (!streamMap.empty())
2571 {
2573 jvObj[jss::type] = "consensusPhase";
2574 jvObj[jss::consensus] = to_string(phase);
2575
2576 for (auto i = streamMap.begin(); i != streamMap.end();)
2577 {
2578 if (auto p = i->second.lock())
2579 {
2580 p->send(jvObj, true);
2581 toRelease.push_back(std::move(p));
2582 ++i;
2583 }
2584 else
2585 {
2586 i = streamMap.erase(i);
2587 }
2588 }
2589 }
2590}
2591
2592void
2594{
2595 // Hold each locked subscriber alive until after streamLock_ is released; a
2596 // last-reference ~InfoSub would otherwise re-acquire this non-recursive
2597 // mutex and self-deadlock. Declared before the lock, destroyed after it.
2599
2600 // VFALCO consider std::shared_mutex
2601 std::scoped_lock const sl(streamLock_);
2602
2603 if (!streamMaps_[SValidations].empty())
2604 {
2606
2607 auto const signerPublic = val->getSignerPublic();
2608
2609 jvObj[jss::type] = "validationReceived";
2610 jvObj[jss::validation_public_key] = toBase58(TokenType::NodePublic, signerPublic);
2611 jvObj[jss::ledger_hash] = to_string(val->getLedgerHash());
2612 jvObj[jss::signature] = strHex(val->getSignature());
2613 jvObj[jss::full] = val->isFull();
2614 jvObj[jss::flags] = val->getFlags();
2615 jvObj[jss::signing_time] = *(*val)[~sfSigningTime];
2616 jvObj[jss::data] = strHex(val->getSerializer().slice());
2617 jvObj[jss::network_id] = registry_.get().getNetworkIDService().getNetworkID();
2618
2619 if (auto version = (*val)[~sfServerVersion])
2620 jvObj[jss::server_version] = std::to_string(*version);
2621
2622 if (auto cookie = (*val)[~sfCookie])
2623 jvObj[jss::cookie] = std::to_string(*cookie);
2624
2625 if (auto hash = (*val)[~sfValidatedHash])
2626 jvObj[jss::validated_hash] = strHex(*hash);
2627
2628 auto const masterKey = registry_.get().getValidatorManifests().getMasterKey(signerPublic);
2629
2630 if (masterKey != signerPublic)
2631 jvObj[jss::master_key] = toBase58(TokenType::NodePublic, masterKey);
2632
2633 // NOTE *seq is a number, but old API versions used string. We replace
2634 // number with a string using MultiApiJson near end of this function
2635 if (auto const seq = (*val)[~sfLedgerSequence])
2636 jvObj[jss::ledger_index] = *seq;
2637
2638 if (val->isFieldPresent(sfAmendments))
2639 {
2640 jvObj[jss::amendments] = json::Value(json::ValueType::Array);
2641 for (auto const& amendment : val->getFieldV256(sfAmendments))
2642 jvObj[jss::amendments].append(to_string(amendment));
2643 }
2644
2645 if (auto const closeTime = (*val)[~sfCloseTime])
2646 jvObj[jss::close_time] = *closeTime;
2647
2648 if (auto const loadFee = (*val)[~sfLoadFee])
2649 jvObj[jss::load_fee] = *loadFee;
2650
2651 if (auto const baseFee = val->at(~sfBaseFee))
2652 jvObj[jss::base_fee] = static_cast<double>(*baseFee);
2653
2654 if (auto const reserveBase = val->at(~sfReserveBase))
2655 jvObj[jss::reserve_base] = *reserveBase;
2656
2657 if (auto const reserveInc = val->at(~sfReserveIncrement))
2658 jvObj[jss::reserve_inc] = *reserveInc;
2659
2660 // (The ~ operator converts the Proxy to a std::optional, which
2661 // simplifies later operations)
2662 if (auto const baseFeeXRP = ~val->at(~sfBaseFeeDrops); baseFeeXRP && baseFeeXRP->native())
2663 jvObj[jss::base_fee] = baseFeeXRP->xrp().jsonClipped();
2664
2665 if (auto const reserveBaseXRP = ~val->at(~sfReserveBaseDrops);
2666 reserveBaseXRP && reserveBaseXRP->native())
2667 jvObj[jss::reserve_base] = reserveBaseXRP->xrp().jsonClipped();
2668
2669 if (auto const reserveIncXRP = ~val->at(~sfReserveIncrementDrops);
2670 reserveIncXRP && reserveIncXRP->native())
2671 jvObj[jss::reserve_inc] = reserveIncXRP->xrp().jsonClipped();
2672
2673 // NOTE Use MultiApiJson to publish two slightly different JSON objects
2674 // for consumers supporting different API versions
2675 MultiApiJson multiObj{jvObj};
2676 multiObj.visit(
2678 [](json::Value& jvTx) {
2679 // Type conversion for older API versions to string
2680 if (jvTx.isMember(jss::ledger_index))
2681 {
2682 jvTx[jss::ledger_index] = std::to_string(jvTx[jss::ledger_index].asUInt());
2683 }
2684 });
2685
2686 for (auto i = streamMaps_[SValidations].begin(); i != streamMaps_[SValidations].end();)
2687 {
2688 if (auto p = i->second.lock())
2689 {
2690 multiObj.visit(
2691 p->getApiVersion(), //
2692 [&](json::Value const& jv) { p->send(jv, true); });
2693 toRelease.push_back(std::move(p));
2694 ++i;
2695 }
2696 else
2697 {
2698 i = streamMaps_[SValidations].erase(i);
2699 }
2700 }
2701 }
2702}
2703
2704void
2706{
2707 // Hold each locked subscriber alive until after streamLock_ is released; a
2708 // last-reference ~InfoSub would otherwise re-acquire this non-recursive
2709 // mutex and self-deadlock. Declared before the lock, destroyed after it.
2711
2712 std::scoped_lock const sl(streamLock_);
2713
2714 if (!streamMaps_[SPeerStatus].empty())
2715 {
2716 json::Value jvObj(func());
2717
2718 jvObj[jss::type] = "peerStatusChange";
2719
2720 for (auto i = streamMaps_[SPeerStatus].begin(); i != streamMaps_[SPeerStatus].end();)
2721 {
2722 InfoSub::pointer p = i->second.lock();
2723
2724 if (p)
2725 {
2726 p->send(jvObj, true);
2727 toRelease.push_back(std::move(p));
2728 ++i;
2729 }
2730 else
2731 {
2732 i = streamMaps_[SPeerStatus].erase(i);
2733 }
2734 }
2735 }
2736}
2737
2738void
2740{
2741 using namespace std::chrono_literals;
2742 if (om == OperatingMode::CONNECTED)
2743 {
2744 if (registry_.get().getLedgerMaster().getValidatedLedgerAge() < 1min)
2746 }
2747 else if (om == OperatingMode::SYNCING)
2748 {
2749 if (registry_.get().getLedgerMaster().getValidatedLedgerAge() >= 1min)
2751 }
2752
2753 if ((om > OperatingMode::CONNECTED) && isBlocked())
2755
2756 if (mode_ == om)
2757 return;
2758
2759 mode_ = om;
2760
2761 accounting_.mode(om);
2762
2763 JLOG(journal_.info()) << "STATE->" << strOperatingMode();
2764 pubServer();
2765}
2766
2767bool
2769{
2770 JLOG(journal_.trace()) << "recvValidation " << val->getLedgerHash() << " from " << source;
2771
2773 BypassAccept bypassAccept = BypassAccept::No;
2774 try
2775 {
2776 if (pendingValidations_.contains(val->getLedgerHash()))
2777 {
2778 bypassAccept = BypassAccept::Yes;
2779 }
2780 else
2781 {
2782 pendingValidations_.insert(val->getLedgerHash());
2783 }
2784 ScopeUnlock const unlock(lock);
2785 handleNewValidation(registry_.get().getApp(), val, source, bypassAccept, journal_);
2786 }
2787 catch (std::exception const& e)
2788 {
2789 JLOG(journal_.warn()) << "Exception thrown for handling new validation "
2790 << val->getLedgerHash() << ": " << e.what();
2791 }
2792 catch (...)
2793 {
2794 JLOG(journal_.warn()) << "Unknown exception thrown for handling new validation "
2795 << val->getLedgerHash();
2796 }
2797 if (bypassAccept == BypassAccept::No)
2798 {
2799 pendingValidations_.erase(val->getLedgerHash());
2800 }
2801 lock.unlock();
2802
2803 pubValidation(val);
2804
2805 JLOG(journal_.debug()) << [this, &val]() -> auto {
2807 ss << "VALIDATION: " << val->render() << " master_key: ";
2808 auto master = registry_.get().getValidators().getTrustedKey(val->getSignerPublic());
2809 if (master)
2810 {
2811 ss << toBase58(TokenType::NodePublic, *master);
2812 }
2813 else
2814 {
2815 ss << "none";
2816 }
2817 return ss.str();
2818 }();
2819
2820 // We will always relay trusted validations; if configured, we will
2821 // also relay all untrusted validations.
2822 return registry_.get().getApp().config().relayUntrustedValidations == 1 || val->isTrusted();
2823}
2824
2827{
2828 return consensus_.getJson(true);
2829}
2830
2832NetworkOPsImp::getServerInfo(bool human, bool admin, bool counters)
2833{
2835
2836 // System-level warnings
2837 {
2839 if (isAmendmentBlocked())
2840 {
2842 w[jss::id] = WarnRpcAmendmentBlocked;
2843 w[jss::message] =
2844 "This server is amendment blocked, and must be updated to be "
2845 "able to stay in sync with the network.";
2846 }
2847 if (isUNLBlocked())
2848 {
2850 w[jss::id] = WarnRpcExpiredValidatorList;
2851 w[jss::message] =
2852 "This server has an expired validator list. validators.txt "
2853 "may be incorrectly configured or some [validator_list_sites] "
2854 "may be unreachable.";
2855 }
2856 if (admin && isAmendmentWarned())
2857 {
2859 w[jss::id] = WarnRpcUnsupportedMajority;
2860 w[jss::message] =
2861 "One or more unsupported amendments have reached majority. "
2862 "Upgrade to the latest version before they are activated "
2863 "to avoid being amendment blocked.";
2864 if (auto const expected =
2865 registry_.get().getAmendmentTable().firstUnsupportedExpected())
2866 {
2867 auto& d = w[jss::details] = json::ValueType::Object;
2868 d[jss::expected_date] = expected->time_since_epoch().count();
2869 d[jss::expected_date_UTC] = to_string(*expected);
2870 }
2871 }
2872
2873 if (warnings.size() != 0u)
2874 info[jss::warnings] = std::move(warnings);
2875 }
2876
2877 // hostid: unique string describing the machine
2878 if (human)
2879 info[jss::hostid] = getHostId(admin);
2880
2881 // domain: if configured with a domain, report it:
2882 if (!registry_.get().getApp().config().serverDomain.empty())
2883 info[jss::server_domain] = registry_.get().getApp().config().serverDomain;
2884
2885 info[jss::build_version] = build_info::getVersionString();
2886
2887 info[jss::server_state] = strOperatingMode(admin);
2888
2889 info[jss::time] =
2890 to_string(std::chrono::floor<std::chrono::microseconds>(std::chrono::system_clock::now()));
2891
2893 info[jss::network_ledger] = "waiting";
2894
2895 info[jss::validation_quorum] =
2896 static_cast<json::UInt>(registry_.get().getValidators().quorum());
2897
2898 if (admin)
2899 {
2900 // Note: By default the node size is "tiny". When parsing it's an error if the final
2901 // NODE_SIZE is over 4 so below code should be safe.
2902 // NOLINTNEXTLINE(bugprone-switch-missing-default-case)
2903 switch (registry_.get().getApp().config().nodeSize)
2904 {
2905 case 0:
2906 info[jss::node_size] = "tiny";
2907 break;
2908 case 1:
2909 info[jss::node_size] = "small";
2910 break;
2911 case 2:
2912 info[jss::node_size] = "medium";
2913 break;
2914 case 3:
2915 info[jss::node_size] = "large";
2916 break;
2917 case 4:
2918 info[jss::node_size] = "huge";
2919 break;
2920 }
2921
2922 auto when = registry_.get().getValidators().expires();
2923
2924 if (!human)
2925 {
2926 if (when)
2927 {
2928 info[jss::validator_list_expires] =
2929 safeCast<json::UInt>(when->time_since_epoch().count());
2930 }
2931 else
2932 {
2933 info[jss::validator_list_expires] = 0;
2934 }
2935 }
2936 else
2937 {
2938 auto& x = (info[jss::validator_list] = json::ValueType::Object);
2939
2940 x[jss::count] = static_cast<json::UInt>(registry_.get().getValidators().count());
2941
2942 if (when)
2943 {
2944 if (*when == TimeKeeper::time_point::max())
2945 {
2946 x[jss::expiration] = "never";
2947 x[jss::status] = "active";
2948 }
2949 else
2950 {
2951 x[jss::expiration] = to_string(*when);
2952
2953 if (*when > registry_.get().getTimeKeeper().now())
2954 {
2955 x[jss::status] = "active";
2956 }
2957 else
2958 {
2959 x[jss::status] = "expired";
2960 }
2961 }
2962 }
2963 else
2964 {
2965 x[jss::status] = "unknown";
2966 x[jss::expiration] = "unknown";
2967 }
2968 }
2969
2970 if (!xrpl::git::getCommitHash().empty() || !xrpl::git::getBuildBranch().empty())
2971 {
2972 auto& x = (info[jss::git] = json::ValueType::Object);
2973 if (!xrpl::git::getCommitHash().empty())
2974 x[jss::hash] = xrpl::git::getCommitHash();
2975 if (!xrpl::git::getBuildBranch().empty())
2976 x[jss::branch] = xrpl::git::getBuildBranch();
2977 }
2978 }
2979 info[jss::io_latency_ms] =
2980 static_cast<json::UInt>(registry_.get().getApp().getIOLatency().count());
2981
2982 if (admin)
2983 {
2984 if (auto const localPubKey = registry_.get().getValidators().localPublicKey();
2985 localPubKey && registry_.get().getApp().getValidationPublicKey())
2986 {
2987 info[jss::pubkey_validator] = toBase58(TokenType::NodePublic, localPubKey.value());
2988 }
2989 else
2990 {
2991 info[jss::pubkey_validator] = "none";
2992 }
2993 }
2994
2995 if (counters)
2996 {
2997 info[jss::counters] = registry_.get().getPerfLog().countersJson();
2998
3000 registry_.get().getNodeStore().getCountsJson(nodestore);
3001 info[jss::counters][jss::nodestore] = nodestore;
3002 info[jss::current_activities] = registry_.get().getPerfLog().currentJson();
3003 }
3004
3005 info[jss::pubkey_node] =
3006 toBase58(TokenType::NodePublic, registry_.get().getApp().nodeIdentity().first);
3007
3008 info[jss::complete_ledgers] = registry_.get().getLedgerMaster().getCompleteLedgers();
3009
3011 info[jss::amendment_blocked] = true;
3012
3013 auto const fp = ledgerMaster_.getFetchPackCacheSize();
3014
3015 if (fp != 0)
3016 info[jss::fetch_pack] = json::UInt(fp);
3017
3018 info[jss::peers] = json::UInt(registry_.get().getOverlay().size());
3019
3021 lastClose[jss::proposers] = json::UInt(consensus_.prevProposers());
3022
3023 if (human)
3024 {
3025 lastClose[jss::converge_time_s] =
3027 }
3028 else
3029 {
3030 lastClose[jss::converge_time] = json::Int(consensus_.prevRoundTime().count());
3031 }
3032
3033 info[jss::last_close] = lastClose;
3034
3035 // info[jss::consensus] = consensus_.getJson();
3036
3037 if (admin)
3038 info[jss::load] = jobQueue_.getJson();
3039
3040 if (auto const netid = registry_.get().getOverlay().networkID())
3041 info[jss::network_id] = static_cast<json::UInt>(*netid);
3042
3043 auto const escalationMetrics =
3044 registry_.get().getTxQ().getMetrics(*registry_.get().getOpenLedger().current());
3045
3046 auto const loadFactorServer = registry_.get().getFeeTrack().getLoadFactor();
3047 auto const loadBaseServer = registry_.get().getFeeTrack().getLoadBase();
3048 /* Scale the escalated fee level to unitless "load factor".
3049 In practice, this just strips the units, but it will continue
3050 to work correctly if either base value ever changes. */
3051 auto const loadFactorFeeEscalation = mulDiv(
3052 escalationMetrics.openLedgerFeeLevel,
3053 loadBaseServer,
3054 escalationMetrics.referenceFeeLevel)
3056
3057 auto const loadFactor =
3058 std::max(safeCast<std::uint64_t>(loadFactorServer), loadFactorFeeEscalation);
3059
3060 if (!human)
3061 {
3062 info[jss::load_base] = loadBaseServer;
3063 info[jss::load_factor] = trunc32(loadFactor);
3064 info[jss::load_factor_server] = loadFactorServer;
3065
3066 /* json::Value doesn't support uint64, so clamp to max
3067 uint32 value. This is mostly theoretical, since there
3068 probably isn't enough extant XRP to drive the factor
3069 that high.
3070 */
3071 info[jss::load_factor_fee_escalation] = escalationMetrics.openLedgerFeeLevel.jsonClipped();
3072 info[jss::load_factor_fee_queue] = escalationMetrics.minProcessingFeeLevel.jsonClipped();
3073 info[jss::load_factor_fee_reference] = escalationMetrics.referenceFeeLevel.jsonClipped();
3074 }
3075 else
3076 {
3077 info[jss::load_factor] = static_cast<double>(loadFactor) / loadBaseServer;
3078
3079 if (loadFactorServer != loadFactor)
3080 info[jss::load_factor_server] = static_cast<double>(loadFactorServer) / loadBaseServer;
3081
3082 if (admin)
3083 {
3084 std::uint32_t fee = registry_.get().getFeeTrack().getLocalFee();
3085 if (fee != loadBaseServer)
3086 info[jss::load_factor_local] = static_cast<double>(fee) / loadBaseServer;
3087 fee = registry_.get().getFeeTrack().getRemoteFee();
3088 if (fee != loadBaseServer)
3089 info[jss::load_factor_net] = static_cast<double>(fee) / loadBaseServer;
3090 fee = registry_.get().getFeeTrack().getClusterFee();
3091 if (fee != loadBaseServer)
3092 info[jss::load_factor_cluster] = static_cast<double>(fee) / loadBaseServer;
3093 }
3094 if (escalationMetrics.openLedgerFeeLevel != escalationMetrics.referenceFeeLevel &&
3095 (admin || loadFactorFeeEscalation != loadFactor))
3096 {
3097 info[jss::load_factor_fee_escalation] =
3098 escalationMetrics.openLedgerFeeLevel.decimalFromReference(
3099 escalationMetrics.referenceFeeLevel);
3100 }
3101 if (escalationMetrics.minProcessingFeeLevel != escalationMetrics.referenceFeeLevel)
3102 {
3103 info[jss::load_factor_fee_queue] =
3104 escalationMetrics.minProcessingFeeLevel.decimalFromReference(
3105 escalationMetrics.referenceFeeLevel);
3106 }
3107 }
3108
3109 bool valid = false;
3110 auto lpClosed = ledgerMaster_.getValidatedLedger();
3111
3112 if (lpClosed)
3113 {
3114 valid = true;
3115 }
3116 else
3117 {
3118 lpClosed = ledgerMaster_.getClosedLedger();
3119 }
3120
3121 if (lpClosed)
3122 {
3123 XRPAmount const baseFee = lpClosed->fees().base;
3125 l[jss::seq] = json::UInt(lpClosed->header().seq);
3126 l[jss::hash] = to_string(lpClosed->header().hash);
3127
3128 if (!human)
3129 {
3130 l[jss::base_fee] = baseFee.jsonClipped();
3131 l[jss::reserve_base] = lpClosed->fees().reserve.jsonClipped();
3132 l[jss::reserve_inc] = lpClosed->fees().increment.jsonClipped();
3133 l[jss::close_time] =
3134 json::Value::UInt(lpClosed->header().closeTime.time_since_epoch().count());
3135 }
3136 else
3137 {
3138 l[jss::base_fee_xrp] = baseFee.decimalXRP();
3139 l[jss::reserve_base_xrp] = lpClosed->fees().reserve.decimalXRP();
3140 l[jss::reserve_inc_xrp] = lpClosed->fees().increment.decimalXRP();
3141
3142 if (auto const closeOffset = registry_.get().getTimeKeeper().closeOffset();
3143 std::abs(closeOffset.count()) >= 60)
3144 l[jss::close_time_offset] = static_cast<std::uint32_t>(closeOffset.count());
3145
3146 static constexpr std::chrono::seconds kHighAgeThreshold{1000000};
3147 if (ledgerMaster_.haveValidated())
3148 {
3149 auto const age = ledgerMaster_.getValidatedLedgerAge();
3150 l[jss::age] = json::UInt(age < kHighAgeThreshold ? age.count() : 0);
3151 }
3152 else
3153 {
3154 auto lCloseTime = lpClosed->header().closeTime;
3155 auto closeTime = registry_.get().getTimeKeeper().closeTime();
3156 if (lCloseTime <= closeTime)
3157 {
3158 using namespace std::chrono_literals;
3159 auto age = closeTime - lCloseTime;
3160 l[jss::age] = json::UInt(age < kHighAgeThreshold ? age.count() : 0);
3161 }
3162 }
3163 }
3164
3165 if (valid)
3166 {
3167 info[jss::validated_ledger] = l;
3168 }
3169 else
3170 {
3171 info[jss::closed_ledger] = l;
3172 }
3173
3174 auto lpPublished = ledgerMaster_.getPublishedLedger();
3175 if (!lpPublished)
3176 {
3177 info[jss::published_ledger] = "none";
3178 }
3179 else if (lpPublished->header().seq != lpClosed->header().seq)
3180 {
3181 info[jss::published_ledger] = lpPublished->header().seq;
3182 }
3183 }
3184
3185 accounting_.json(info);
3186 info[jss::uptime] = UptimeClock::now().time_since_epoch().count();
3187 info[jss::jq_trans_overflow] =
3188 std::to_string(registry_.get().getOverlay().getJqTransOverflow());
3189 info[jss::peer_disconnects] = std::to_string(registry_.get().getOverlay().getPeerDisconnect());
3190 info[jss::peer_disconnects_resources] =
3191 std::to_string(registry_.get().getOverlay().getPeerDisconnectCharges());
3192
3193 // This array must be sorted in increasing order.
3194 static constexpr std::array<std::string_view, 7> kProtocols{
3195 "http", "https", "peer", "ws", "ws2", "wss", "wss2"};
3196 static_assert(std::ranges::is_sorted(kProtocols));
3197 {
3199 for (auto const& port : registry_.get().getServerHandler().setup().ports)
3200 {
3201 // Don't publish admin ports for non-admin users
3202 if (!admin &&
3203 !(port.adminNetsV4.empty() && port.adminNetsV6.empty() && port.adminUser.empty() &&
3204 port.adminPassword.empty()))
3205 continue;
3207 // NOLINTNEXTLINE(modernize-use-ranges)
3209 std::begin(port.protocol),
3210 std::end(port.protocol),
3211 std::begin(kProtocols),
3212 std::end(kProtocols),
3213 std::back_inserter(proto));
3214 if (!proto.empty())
3215 {
3216 auto& jv = ports.append(json::Value(json::ValueType::Object));
3217 jv[jss::port] = std::to_string(port.port);
3218 jv[jss::protocol] = json::Value{json::ValueType::Array};
3219 for (auto const& p : proto)
3220 jv[jss::protocol].append(p);
3221 }
3222 }
3223
3224 if (registry_.get().getApp().config().exists(Sections::kPortGrpc))
3225 {
3226 auto const& grpcSection =
3227 registry_.get().getApp().config().section(Sections::kPortGrpc);
3228 auto const optPort = grpcSection.get(Keys::kPort);
3229 if (optPort && grpcSection.get(Keys::kIp))
3230 {
3231 auto& jv = ports.append(json::Value(json::ValueType::Object));
3232 jv[jss::port] = *optPort;
3233 jv[jss::protocol] = json::Value{json::ValueType::Array};
3234 jv[jss::protocol].append("grpc");
3235 }
3236 }
3237 info[jss::ports] = std::move(ports);
3238 }
3239
3240 return info;
3241}
3242
3243void
3245{
3246 registry_.get().getInboundLedgers().clearFailures();
3247}
3248
3251{
3252 return registry_.get().getInboundLedgers().getInfo();
3253}
3254
3255void
3257 std::shared_ptr<ReadView const> const& ledger,
3258 std::shared_ptr<STTx const> const& transaction,
3259 TER result)
3260{
3261 // never publish an inner txn inside a batch txn. The flag should
3262 // only be set if the Batch feature is enabled. If Batch is not
3263 // enabled, the flag is always invalid, so don't publish it
3264 // regardless.
3265 if (transaction->isFlag(tfInnerBatchTxn))
3266 return;
3267
3268 MultiApiJson const jvObj = transJson(transaction, result, false, ledger, std::nullopt);
3269
3270 {
3271 // Hold each locked subscriber alive until after streamLock_ is
3272 // released; a last-reference ~InfoSub would otherwise re-acquire this
3273 // non-recursive mutex and self-deadlock. Declared before the lock,
3274 // destroyed after the block ends.
3276
3277 std::scoped_lock const sl(streamLock_);
3278
3279 auto it = streamMaps_[SRtTransactions].begin();
3280 while (it != streamMaps_[SRtTransactions].end())
3281 {
3282 InfoSub::pointer p = it->second.lock();
3283
3284 if (p)
3285 {
3286 jvObj.visit(
3287 p->getApiVersion(), //
3288 [&](json::Value const& jv) { p->send(jv, true); });
3289 toRelease.push_back(std::move(p));
3290 ++it;
3291 }
3292 else
3293 {
3294 it = streamMaps_[SRtTransactions].erase(it);
3295 }
3296 }
3297 }
3298
3299 pubProposedAccountTransaction(ledger, transaction, result);
3300}
3301
3302void
3304{
3305 // Ledgers are published only when they acquire sufficient validations
3306 // Holes are filled across connection loss or other catastrophe
3307
3309 registry_.get().getAcceptedLedgerCache().fetch(lpAccepted->header().hash);
3310 if (!alpAccepted)
3311 {
3312 alpAccepted = std::make_shared<AcceptedLedger>(lpAccepted);
3313 registry_.get().getAcceptedLedgerCache().canonicalizeReplaceClient(
3314 lpAccepted->header().hash, alpAccepted);
3315 }
3316
3317 XRPL_ASSERT(
3318 alpAccepted->getLedger().get() == lpAccepted.get(),
3319 "xrpl::NetworkOPsImp::pubLedger : accepted input");
3320
3321 JLOG(journal_.debug()) << "Publishing ledger " << lpAccepted->header().seq << " "
3322 << lpAccepted->header().hash;
3323
3324 // Stream updates and the account-history kick-off touch different lock
3325 // domains; each helper takes only its own lock, so the two are never held
3326 // together.
3327 publishLedgerStreams(lpAccepted, alpAccepted);
3328 kickoffAccountHistory(alpAccepted);
3329
3330 // Don't lock since pubAcceptedTransaction is locking.
3331 for (auto const& accTx : *alpAccepted)
3332 {
3333 JLOG(journal_.trace()) << "pubAccepted: " << accTx->getJson();
3334 bool const last = &*accTx == &alpAccepted->back();
3335 pubValidatedTransaction(lpAccepted, *accTx, last);
3336 }
3337}
3338
3339void
3341 std::shared_ptr<ReadView const> const& lpAccepted,
3342 std::shared_ptr<AcceptedLedger const> const& alpAccepted)
3343{
3344 // Hold each locked subscriber alive until after streamLock_ is released; a
3345 // last-reference ~InfoSub would otherwise re-acquire this non-recursive
3346 // mutex and self-deadlock. Declared before the lock, destroyed after it;
3347 // covers both the ledger and book-changes loops below.
3349
3350 std::scoped_lock const sl(streamLock_);
3351
3352 if (!streamMaps_[SLedger].empty())
3353 {
3355
3356 jvObj[jss::type] = "ledgerClosed";
3357 jvObj[jss::ledger_index] = lpAccepted->header().seq;
3358 jvObj[jss::ledger_hash] = to_string(lpAccepted->header().hash);
3359 jvObj[jss::ledger_time] =
3360 json::Value::UInt(lpAccepted->header().closeTime.time_since_epoch().count());
3361
3362 jvObj[jss::network_id] = registry_.get().getNetworkIDService().getNetworkID();
3363
3364 if (!lpAccepted->rules().enabled(featureXRPFees))
3365 jvObj[jss::fee_ref] = kFeeUnitsDeprecated;
3366 jvObj[jss::fee_base] = lpAccepted->fees().base.jsonClipped();
3367 jvObj[jss::reserve_base] = lpAccepted->fees().reserve.jsonClipped();
3368 jvObj[jss::reserve_inc] = lpAccepted->fees().increment.jsonClipped();
3369
3370 jvObj[jss::txn_count] = json::UInt(alpAccepted->size());
3371
3373 {
3374 jvObj[jss::validated_ledgers] = registry_.get().getLedgerMaster().getCompleteLedgers();
3375 }
3376 auto it = streamMaps_[SLedger].begin();
3377 while (it != streamMaps_[SLedger].end())
3378 {
3379 InfoSub::pointer p = it->second.lock();
3380 if (p)
3381 {
3382 p->send(jvObj, true);
3383 toRelease.push_back(std::move(p));
3384 ++it;
3385 }
3386 else
3387 {
3388 it = streamMaps_[SLedger].erase(it);
3389 }
3390 }
3391 }
3392
3393 if (!streamMaps_[SBookChanges].empty())
3394 {
3395 json::Value const jvObj = xrpl::rpc::computeBookChanges(lpAccepted);
3396
3397 auto it = streamMaps_[SBookChanges].begin();
3398 while (it != streamMaps_[SBookChanges].end())
3399 {
3400 InfoSub::pointer p = it->second.lock();
3401 if (p)
3402 {
3403 p->send(jvObj, true);
3404 toRelease.push_back(std::move(p));
3405 ++it;
3406 }
3407 else
3408 {
3409 it = streamMaps_[SBookChanges].erase(it);
3410 }
3411 }
3412 }
3413}
3414
3415void
3417{
3418 // Runs exactly once, the first time a ledger is published. The atomic
3419 // exchange lets the common post-first-ledger path return without taking
3420 // accountLock_, while still admitting exactly one caller even if ledger
3421 // publishing is ever made concurrent.
3422 static std::atomic<bool> done{false};
3423 if (done.exchange(true))
3424 return;
3425
3426 // It only reads/writes subAccountHistory_, so it takes accountLock_ alone.
3428 for (auto& outer : subAccountHistory_)
3429 {
3430 for (auto& inner : outer.second)
3431 {
3432 auto& subInfo = inner.second;
3433 if (subInfo.index->separationLedgerSeq == 0)
3434 subAccountHistoryStart(alpAccepted->getLedger(), subInfo);
3435 }
3436 }
3437}
3438
3439void
3441{
3442 ServerFeeSummary const f{
3443 registry_.get().getOpenLedger().current()->fees().base,
3444 registry_.get().getTxQ().getMetrics(*registry_.get().getOpenLedger().current()),
3445 registry_.get().getFeeTrack()};
3446
3447 // only schedule the job if something has changed
3448 if (f != lastFeeSummary_)
3449 {
3450 jobQueue_.addJob(JtClientFeeChange, "PubFee", [this]() { pubServer(); });
3451 }
3452}
3453
3454void
3456{
3457 jobQueue_.addJob(JtClientConsensus, "PubCons", [this, phase]() { pubConsensus(phase); });
3458}
3459
3460inline void
3462{
3463 localTX_->sweep(view);
3464}
3465inline std::size_t
3467{
3468 return localTX_->size();
3469}
3470
3473{
3474 std::scoped_lock const sl(bookLock_);
3475 std::size_t total = 0;
3476 for (auto const& [_, subs] : subBook_)
3477 total += subs.size();
3478 return total;
3479}
3480
3481// This routine should only be used to publish accepted or validated
3482// transactions.
3485 std::shared_ptr<STTx const> const& transaction,
3486 TER result,
3487 bool validated,
3488 std::shared_ptr<ReadView const> const& ledger,
3490{
3492 std::string sToken;
3493 std::string sHuman;
3494
3495 transResultInfo(result, sToken, sHuman);
3496
3497 jvObj[jss::type] = "transaction";
3498 // NOTE jvObj is not a finished object for either API version. After
3499 // it's populated, we need to finish it for a specific API version. This is
3500 // done in a loop, near the end of this function.
3501 jvObj[jss::transaction] = transaction->getJson(JsonOptions::Values::DisableApiPriorV2, false);
3502
3503 if (meta)
3504 {
3505 jvObj[jss::meta] = meta->get().getJson(JsonOptions::Values::None);
3506 rpc::insertAllSyntheticInJson(jvObj[jss::meta], *ledger, transaction, meta->get());
3507 }
3508
3509 // add CTID where the needed data for it exists
3510 if (auto const& lookup = ledger->txRead(transaction->getTransactionID());
3511 lookup.second && lookup.second->isFieldPresent(sfTransactionIndex))
3512 {
3513 uint32_t const txnSeq = lookup.second->getFieldU32(sfTransactionIndex);
3514 uint32_t netID = registry_.get().getNetworkIDService().getNetworkID();
3515 if (transaction->isFieldPresent(sfNetworkID))
3516 netID = transaction->getFieldU32(sfNetworkID);
3517
3518 if (std::optional<std::string> ctid = rpc::encodeCTID(ledger->header().seq, txnSeq, netID);
3519 ctid)
3520 jvObj[jss::ctid] = *ctid;
3521 }
3522 if (!ledger->open())
3523 jvObj[jss::ledger_hash] = to_string(ledger->header().hash);
3524
3525 if (validated)
3526 {
3527 jvObj[jss::ledger_index] = ledger->header().seq;
3528 jvObj[jss::transaction][jss::date] = ledger->header().closeTime.time_since_epoch().count();
3529 jvObj[jss::validated] = true;
3530 jvObj[jss::close_time_iso] = toStringIso(ledger->header().closeTime);
3531
3532 // WRITEME: Put the account next seq here
3533 }
3534 else
3535 {
3536 jvObj[jss::validated] = false;
3537 jvObj[jss::ledger_current_index] = ledger->header().seq;
3538 }
3539
3540 jvObj[jss::status] = validated ? "closed" : "proposed";
3541 jvObj[jss::engine_result] = sToken;
3542 jvObj[jss::engine_result_code] = result;
3543 jvObj[jss::engine_result_message] = sHuman;
3544
3545 if (transaction->getTxnType() == ttOFFER_CREATE)
3546 {
3547 auto const account = transaction->getAccountID(sfAccount);
3548 auto const amount = transaction->getFieldAmount(sfTakerGets);
3549
3550 // If the offer create is not self funded then add the owner balance
3551 if (account != amount.getIssuer())
3552 {
3553 auto const ownerFunds = accountFunds(
3554 *ledger,
3555 account,
3556 amount,
3559 registry_.get().getJournal("View"));
3560 jvObj[jss::transaction][jss::owner_funds] = ownerFunds.getText();
3561 }
3562 }
3563
3564 std::string const hash = to_string(transaction->getTransactionID());
3565 MultiApiJson multiObj{jvObj};
3567 multiObj.visit(), //
3568 [&]<unsigned Version>(json::Value& jvTx, std::integral_constant<unsigned, Version>) {
3569 rpc::insertDeliverMax(jvTx[jss::transaction], transaction->getTxnType(), Version);
3570
3571 if constexpr (Version > 1)
3572 {
3573 jvTx[jss::tx_json] = jvTx.removeMember(jss::transaction);
3574 jvTx[jss::hash] = hash;
3575 }
3576 else
3577 {
3578 jvTx[jss::transaction][jss::hash] = hash;
3579 }
3580 });
3581
3582 return multiObj;
3583}
3584
3585void
3587 std::shared_ptr<ReadView const> const& ledger,
3588 AcceptedLedgerTx const& transaction,
3589 bool last)
3590{
3591 auto const& stTxn = transaction.getTxn();
3592
3593 // Create two different Json objects, for different API versions
3594 auto const metaRef = std::ref(transaction.getMeta());
3595 auto const trResult = transaction.getResult();
3596 MultiApiJson const jvObj = transJson(stTxn, trResult, true, ledger, metaRef);
3597
3598 {
3599 // Hold each locked subscriber alive until after streamLock_ is
3600 // released; a last-reference ~InfoSub would otherwise re-acquire this
3601 // non-recursive mutex and self-deadlock. Declared before the lock,
3602 // destroyed after the block ends; covers both loops below.
3604
3605 std::scoped_lock const sl(streamLock_);
3606
3607 auto it = streamMaps_[STransactions].begin();
3608 while (it != streamMaps_[STransactions].end())
3609 {
3610 InfoSub::pointer p = it->second.lock();
3611
3612 if (p)
3613 {
3614 jvObj.visit(
3615 p->getApiVersion(), //
3616 [&](json::Value const& jv) { p->send(jv, true); });
3617 toRelease.push_back(std::move(p));
3618 ++it;
3619 }
3620 else
3621 {
3622 it = streamMaps_[STransactions].erase(it);
3623 }
3624 }
3625
3626 it = streamMaps_[SRtTransactions].begin();
3627
3628 while (it != streamMaps_[SRtTransactions].end())
3629 {
3630 InfoSub::pointer p = it->second.lock();
3631
3632 if (p)
3633 {
3634 jvObj.visit(
3635 p->getApiVersion(), //
3636 [&](json::Value const& jv) { p->send(jv, true); });
3637 toRelease.push_back(std::move(p));
3638 ++it;
3639 }
3640 else
3641 {
3642 it = streamMaps_[SRtTransactions].erase(it);
3643 }
3644 }
3645 }
3646
3647 if (transaction.getResult() == tesSUCCESS)
3648 pubBookTransaction(transaction, jvObj);
3649
3650 pubAccountTransaction(ledger, transaction, last);
3651 pubMPTTransaction(transaction, jvObj);
3652}
3653
3654void
3656{
3657 auto const books = affectedBooks(alTx, journal_);
3658 if (books.empty())
3659 return;
3660
3661 // Two-pass design:
3662 //
3663 // 1. Under bookLock_, walk subBook_, collect a strong pointer for each
3664 // unique listener (and prune any expired weak_ptrs we encounter).
3665 // 2. Release bookLock_, then send to each collected listener.
3666 //
3667 // Reasoning:
3668 // * send() can be slow / blocking, so holding bookLock_ across it would
3669 // stall every other book sub/unsub/pub path on this server (see the
3670 // matching TODO above pubServer at line ~2275).
3671 // * A strong pointer destructed while bookLock_ is held risks running
3672 // ~InfoSub() in-line, which re-enters unsubBook() and mutates the very
3673 // subBook_/SubMapType being iterated -> dangling iterator UB.
3674 //
3675 // Releasing bookLock_ before any InfoSub::pointer can decay solves both.
3676 // ~InfoSub() reacquires bookLock_ via unsubBook() on its own and serializes
3677 // safely with concurrent traffic.
3678
3681
3682 // Sized for the common case where every affected book has at most
3683 // one subscriber. Multi-subscriber books trigger reallocation, but
3684 // that is rare and the upper-bound estimate (sum of per-book sizes)
3685 // would itself require walking subBook_ twice.
3686 listeners.reserve(books.size());
3687 seen.reserve(books.size());
3688
3689 {
3690 std::scoped_lock const sl(bookLock_);
3691
3692 for (auto const& book : books)
3693 {
3694 auto it = subBook_.find(book);
3695 if (it == subBook_.end())
3696 continue;
3697
3698 for (auto sit = it->second.begin(); sit != it->second.end();)
3699 {
3700 if (auto p = sit->second.lock())
3701 {
3702 // Defensive: subBook_ entries are normally cleared by
3703 // ~InfoSub() -> unsubBook(), so we rarely see expired
3704 // weak_ptrs here. The else branch covers the narrow race
3705 // where the last strong ref is dropped between insertion
3706 // and our lock() call.
3707 if (seen.emplace(p->getSeq()).second)
3708 listeners.emplace_back(std::move(p));
3709 ++sit;
3710 }
3711 else
3712 {
3713 JLOG(journal_.debug())
3714 << "pubBookTransaction: pruning expired weak_ptr for seq=" << sit->first;
3715 sit = it->second.erase(sit);
3716 }
3717 }
3718
3719 if (it->second.empty())
3720 subBook_.erase(it);
3721 }
3722 }
3723
3724 for (auto const& p : listeners)
3725 {
3726 jvObj.visit(p->getApiVersion(), [&](json::Value const& jv) { p->send(jv, true); });
3727 }
3728 // listeners destructs here, outside bookLock_; ~InfoSub (if any fires)
3729 // will reacquire bookLock_ via unsubBook with no iterator hazard.
3730}
3731
3732void
3734 std::shared_ptr<ReadView const> const& ledger,
3735 AcceptedLedgerTx const& transaction,
3736 bool last)
3737{
3739 int iProposed = 0;
3740 int iAccepted = 0;
3741
3742 std::vector<SubAccountHistoryInfo> accountHistoryNotify;
3743 auto const currLedgerSeq = ledger->seq();
3744 {
3746
3747 if (!subAccount_.empty() || !subRTAccount_.empty() || !subAccountHistory_.empty())
3748 {
3749 for (auto const& affectedAccount : transaction.getAffected())
3750 {
3751 if (auto simiIt = subRTAccount_.find(affectedAccount);
3752 simiIt != subRTAccount_.end())
3753 {
3754 auto it = simiIt->second.begin();
3755
3756 while (it != simiIt->second.end())
3757 {
3758 InfoSub::pointer const p = it->second.lock();
3759
3760 if (p)
3761 {
3762 notify.insert(p);
3763 ++it;
3764 ++iProposed;
3765 }
3766 else
3767 {
3768 it = simiIt->second.erase(it);
3769 }
3770 }
3771 }
3772
3773 if (auto simiIt = subAccount_.find(affectedAccount); simiIt != subAccount_.end())
3774 {
3775 auto it = simiIt->second.begin();
3776 while (it != simiIt->second.end())
3777 {
3778 InfoSub::pointer const p = it->second.lock();
3779
3780 if (p)
3781 {
3782 notify.insert(p);
3783 ++it;
3784 ++iAccepted;
3785 }
3786 else
3787 {
3788 it = simiIt->second.erase(it);
3789 }
3790 }
3791 }
3792
3793 if (auto historyIt = subAccountHistory_.find(affectedAccount);
3794 historyIt != subAccountHistory_.end())
3795 {
3796 auto& subs = historyIt->second;
3797 auto it = subs.begin();
3798 while (it != subs.end())
3799 {
3800 SubAccountHistoryInfoWeak const& info = it->second;
3801 if (currLedgerSeq <= info.index->separationLedgerSeq)
3802 {
3803 ++it;
3804 continue;
3805 }
3806
3807 if (auto isSptr = info.sinkWptr.lock(); isSptr)
3808 {
3809 accountHistoryNotify.emplace_back(
3810 SubAccountHistoryInfo{.sink = isSptr, .index = info.index});
3811 ++it;
3812 }
3813 else
3814 {
3815 it = subs.erase(it);
3816 }
3817 }
3818 if (subs.empty())
3819 subAccountHistory_.erase(historyIt);
3820 }
3821 }
3822 }
3823 }
3824
3825 JLOG(journal_.trace()) << "pubAccountTransaction: "
3826 << "proposed=" << iProposed << ", accepted=" << iAccepted;
3827
3828 if (!notify.empty() || !accountHistoryNotify.empty())
3829 {
3830 auto const& stTxn = transaction.getTxn();
3831
3832 // Create two different Json objects, for different API versions
3833 auto const metaRef = std::ref(transaction.getMeta());
3834 auto const trResult = transaction.getResult();
3835 MultiApiJson jvObj = transJson(stTxn, trResult, true, ledger, metaRef);
3836
3837 for (InfoSub::Ref isrListener : notify)
3838 {
3839 jvObj.visit(
3840 isrListener->getApiVersion(), //
3841 [&](json::Value const& jv) { isrListener->send(jv, true); });
3842 }
3843
3844 if (last)
3845 jvObj.set(jss::account_history_boundary, true);
3846
3847 XRPL_ASSERT(
3848 jvObj.isMember(jss::account_history_tx_stream) == MultiApiJson::IsMemberResult::None,
3849 "xrpl::NetworkOPsImp::pubAccountTransaction : "
3850 "account_history_tx_stream not set");
3851 for (auto& info : accountHistoryNotify)
3852 {
3853 auto& index = info.index;
3854 if (index->forwardTxIndex == 0 && !index->haveHistorical)
3855 jvObj.set(jss::account_history_tx_first, true);
3856
3857 jvObj.set(jss::account_history_tx_index, index->forwardTxIndex++);
3858
3859 jvObj.visit(
3860 info.sink->getApiVersion(), //
3861 [&](json::Value const& jv) { info.sink->send(jv, true); });
3862 }
3863 }
3864}
3865
3866void
3868 std::shared_ptr<ReadView const> const& ledger,
3870 TER result)
3871{
3873 int iProposed = 0;
3874
3875 std::vector<SubAccountHistoryInfo> accountHistoryNotify;
3876
3877 {
3879
3880 if (subRTAccount_.empty())
3881 return;
3882
3883 if (!subAccount_.empty() || !subRTAccount_.empty() || !subAccountHistory_.empty())
3884 {
3885 for (auto const& affectedAccount : tx->getMentionedAccounts())
3886 {
3887 if (auto simiIt = subRTAccount_.find(affectedAccount);
3888 simiIt != subRTAccount_.end())
3889 {
3890 auto it = simiIt->second.begin();
3891
3892 while (it != simiIt->second.end())
3893 {
3894 InfoSub::pointer const p = it->second.lock();
3895
3896 if (p)
3897 {
3898 notify.insert(p);
3899 ++it;
3900 ++iProposed;
3901 }
3902 else
3903 {
3904 it = simiIt->second.erase(it);
3905 }
3906 }
3907 }
3908 }
3909 }
3910 }
3911
3912 JLOG(journal_.trace()) << "pubProposedAccountTransaction: " << iProposed;
3913
3914 if (!notify.empty() || !accountHistoryNotify.empty())
3915 {
3916 // Create two different Json objects, for different API versions
3917 MultiApiJson jvObj = transJson(tx, result, false, ledger, std::nullopt);
3918
3919 for (InfoSub::Ref isrListener : notify)
3920 {
3921 jvObj.visit(
3922 isrListener->getApiVersion(), //
3923 [&](json::Value const& jv) { isrListener->send(jv, true); });
3924 }
3925
3926 XRPL_ASSERT(
3927 jvObj.isMember(jss::account_history_tx_stream) == MultiApiJson::IsMemberResult::None,
3928 "xrpl::NetworkOPs::pubProposedAccountTransaction : "
3929 "account_history_tx_stream not set");
3930 for (auto& info : accountHistoryNotify)
3931 {
3932 auto& index = info.index;
3933 if (index->forwardTxIndex == 0 && !index->haveHistorical)
3934 jvObj.set(jss::account_history_tx_first, true);
3935 jvObj.set(jss::account_history_tx_index, index->forwardTxIndex++);
3936 jvObj.visit(
3937 info.sink->getApiVersion(), //
3938 [&](json::Value const& jv) { info.sink->send(jv, true); });
3939 }
3940 }
3941}
3942
3943//
3944// Monitoring
3945//
3946
3947void
3949 InfoSub::Ref isrListener,
3950 HashSet<AccountID> const& vnaAccountIDs,
3951 bool rt)
3952{
3953 SubInfoMapType& subMap = rt ? subRTAccount_ : subAccount_;
3954
3955 for (auto const& naAccountID : vnaAccountIDs)
3956 {
3957 JLOG(journal_.trace()) << "subAccount: account: " << toBase58(naAccountID);
3958
3959 isrListener->insertSubAccountInfo(naAccountID, rt);
3960 }
3961
3963
3964 for (auto const& naAccountID : vnaAccountIDs)
3965 {
3966 auto simIterator = subMap.find(naAccountID);
3967 if (simIterator == subMap.end())
3968 {
3969 // Not found, note that account has a new single listener.
3970 SubMapType usisElement;
3971 usisElement[isrListener->getSeq()] = isrListener;
3972 // VFALCO NOTE This is making a needless copy of naAccountID
3973 subMap.insert(simIterator, make_pair(naAccountID, usisElement));
3974 }
3975 else
3976 {
3977 // Found, note that the account has another listener.
3978 simIterator->second[isrListener->getSeq()] = isrListener;
3979 }
3980 }
3981}
3982
3983void
3985 InfoSub::Ref isrListener,
3986 HashSet<AccountID> const& vnaAccountIDs,
3987 bool rt)
3988{
3989 for (auto const& naAccountID : vnaAccountIDs)
3990 {
3991 // Remove from the InfoSub
3992 isrListener->deleteSubAccountInfo(naAccountID, rt);
3993 }
3994
3995 // Remove from the server
3996 unsubAccountInternal(isrListener->getSeq(), vnaAccountIDs, rt);
3997}
3998
3999void
4001 std::uint64_t uSeq,
4002 HashSet<AccountID> const& vnaAccountIDs,
4003 bool rt)
4004{
4006
4007 SubInfoMapType& subMap = rt ? subRTAccount_ : subAccount_;
4008
4009 for (auto const& naAccountID : vnaAccountIDs)
4010 {
4011 auto simIterator = subMap.find(naAccountID);
4012
4013 if (simIterator != subMap.end())
4014 {
4015 // Found
4016 simIterator->second.erase(uSeq);
4017
4018 if (simIterator->second.empty())
4019 {
4020 // Don't need hash entry.
4021 subMap.erase(simIterator);
4022 }
4023 }
4024 }
4025}
4026
4027template <typename OuterMap, typename BeforeErase>
4028void
4030 std::uint64_t seq,
4031 HashSet<AccountID> const& accounts,
4032 OuterMap& outerMap,
4033 BeforeErase&& beforeErase)
4034{
4035 // Walk the disconnecting connection's accounts in chunks. Each chunk takes
4036 // accountLock_, erases up to kAccountCleanupChunk entries, then releases
4037 // the lock so a competing account-publish can run before the next chunk.
4038 // No iterator into outerMap is held across the unlock: every chunk re-finds
4039 // each account, so a concurrent mutation between chunks cannot dangle.
4040 auto it = accounts.begin();
4041 auto const end = accounts.end();
4042 while (it != end)
4043 {
4045
4046 for (std::size_t n = 0; n < kAccountCleanupChunk && it != end; ++n, ++it)
4047 {
4048 auto outerIter = outerMap.find(*it);
4049 if (outerIter != outerMap.end())
4050 {
4051 // Give the caller a chance to tear down this connection's inner
4052 // entry before it is erased (the history map stops its paging
4053 // job here); the plain account maps pass a no-op.
4054 auto innerIter = outerIter->second.find(seq);
4055 if (innerIter != outerIter->second.end())
4056 beforeErase(innerIter->second);
4057
4058 // Erase only this connection's seq; other connections sharing
4059 // the account keep their entry, so a reconnect is unaffected.
4060 outerIter->second.erase(seq);
4061 if (outerIter->second.empty())
4062 outerMap.erase(outerIter);
4063 }
4064 }
4065 }
4066}
4067
4068void
4070 std::uint64_t seq,
4071 HashSet<AccountID> const& accounts,
4072 SubInfoMapType& subMap)
4073{
4074 // Plain account maps need no per-entry teardown before erase.
4075 cleanupSubscriptionMap(seq, accounts, subMap, [](InfoSub::Wptr const&) {});
4076}
4077
4078void
4080 std::uint64_t seq,
4081 HashSet<AccountID> const& accounts)
4082{
4083 // Cancel any in-flight historical paging job for this connection before
4084 // dropping its record. The job holds its own shared_ptr to the index, so
4085 // erasing the map entry alone would not stop it; it reads this atomic
4086 // between pages and exits promptly once set.
4088 seq, accounts, subAccountHistory_, [](SubAccountHistoryInfoWeak const& info) {
4089 info.index->stopHistorical = true;
4090 });
4091}
4092
4093void
4095 std::uint64_t seq,
4096 HashSet<AccountID> rtAccounts,
4097 HashSet<AccountID> normalAccounts,
4098 HashSet<AccountID> historyAccounts)
4099{
4100 // Nothing to do for a connection that never subscribed to any account.
4101 if (rtAccounts.empty() && normalAccounts.empty() && historyAccounts.empty())
4102 return;
4103
4104 // Post the erase work to a low-priority job so the disconnect thread (and
4105 // ~InfoSub) returns immediately. The job captures the sets BY MOVE and
4106 // operates purely on seq + the captured accounts; it never touches the
4107 // destroyed InfoSub. `this` outlives the job per the Source lifetime
4108 // contract. Running on a JobQueue thread, it cannot re-enter accountLock_
4109 // held by the disconnecting thread, so the plain std::mutex is safe.
4110 //
4111 // The body is exception-guarded: the JobQueue invokes it bare, so an
4112 // escaping exception on the worker thread would terminate the process.
4113 //
4114 // addJob returns false only once the JobQueue has been stopped, i.e. during
4115 // process shutdown. At that point NetworkOPsImp's maps are about to be
4116 // destroyed wholesale and no publish path can run, so dropping the cleanup
4117 // is harmless; no inline fallback is needed.
4118 jobQueue_.addJob(
4120 "SubCleanup",
4121 [this,
4122 seq,
4123 rt = std::move(rtAccounts),
4124 normal = std::move(normalAccounts),
4125 history = std::move(historyAccounts)]() noexcept {
4126 try
4127 {
4131 }
4132 catch (std::exception const& e)
4133 {
4134 JLOG(journal_.error()) << "SubCleanup[seq=" << seq << "]: " << e.what();
4135 }
4136 catch (...)
4137 {
4138 JLOG(journal_.error()) << "SubCleanup[seq=" << seq << "]: unknown exception";
4139 }
4140 });
4141}
4142
4143void
4145{
4146 {
4147 std::scoped_lock const sl(mptLock_);
4148 if (subMPT_.empty())
4149 return;
4150 }
4151
4152 // getAffectedMPTs derives the issuance id for MPTokenIssuance entries from
4153 // the metadata, so it also covers transactions whose top-level STTx carries
4154 // no MPTokenIssuanceID, such as the inner transactions of a Batch. It reads
4155 // only the metadata, so it runs outside mptLock_.
4156 auto const affectedMPTs = alTx.getMeta().getAffectedMPTs();
4157 if (affectedMPTs.empty())
4158 return;
4159
4160 // Declared before the lock so a last-reference ~InfoSub runs after
4161 // mptLock_ is released (see the deferred-destruction rule).
4163
4164 {
4165 std::scoped_lock const sl(mptLock_);
4166
4167 for (auto const& affectedMPT : affectedMPTs)
4168 {
4169 if (auto simiIt = subMPT_.find(affectedMPT); simiIt != subMPT_.end())
4170 {
4171 auto it = simiIt->second.begin();
4172 while (it != simiIt->second.end())
4173 {
4174 InfoSub::pointer const p = it->second.lock();
4175
4176 if (p)
4177 {
4178 notify.insert(p);
4179 ++it;
4180 }
4181 else
4182 {
4183 it = simiIt->second.erase(it);
4184 }
4185 }
4186 }
4187 }
4188 }
4189
4190 if (notify.empty())
4191 return;
4192
4193 // Reuse the transaction JSON built by pubValidatedTransaction; only the
4194 // message type differs for this stream.
4195 MultiApiJson jvMPT = jvObj;
4196 jvMPT.set(jss::type, "mptTransaction");
4197
4198 for (InfoSub::Ref isrListener : notify)
4199 {
4200 jvMPT.visit(isrListener->getApiVersion(), [&](json::Value const& jv) {
4201 isrListener->send(jv, true);
4202 });
4203 }
4204}
4205
4206void
4208{
4209 // Insert into subMPT_ before the InfoSub (unsubMPT removes in reverse), as
4210 // subBook does, so a racing unsubMPT cannot leave a subMPT_ entry that the
4211 // InfoSub doesn't know about.
4212 {
4213 std::scoped_lock const sl(mptLock_);
4214
4215 for (auto const& mptID : mptIDs)
4216 {
4217 JLOG(journal_.trace()) << "subMPT: MPT: " << to_string(mptID);
4218
4219 auto simIterator = subMPT_.find(mptID);
4220 if (simIterator == subMPT_.end())
4221 {
4222 // Not found, note that the MPT issuance has a new single listener.
4223 SubMapType usisElement;
4224 usisElement[isrListener->getSeq()] = isrListener;
4225 subMPT_.insert(simIterator, make_pair(mptID, usisElement));
4226 }
4227 else
4228 {
4229 // Found, note that the MPT issuance has another listener.
4230 simIterator->second[isrListener->getSeq()] = isrListener;
4231 }
4232 }
4233 }
4234
4235 for (auto const& mptID : mptIDs)
4236 isrListener->insertSubMPTInfo(mptID);
4237}
4238
4239void
4241{
4242 for (auto const& mptID : mptIDs)
4243 {
4244 // Remove from the InfoSub first so ~InfoSub does not re-issue an
4245 // unsubMPTInternal for an issuance the caller already removed, then
4246 // remove from the server.
4247 isrListener->deleteSubMPTInfo(mptID);
4248 unsubMPTInternal(isrListener->getSeq(), mptID);
4249 }
4250}
4251
4252void
4254{
4255 // Only weak_ptrs are erased, so no InfoSub can be destroyed under the lock.
4256 std::scoped_lock const sl(mptLock_);
4257
4258 auto simIterator = subMPT_.find(mptID);
4259 if (simIterator == subMPT_.end())
4260 return;
4261
4262 simIterator->second.erase(uSeq);
4263 if (simIterator->second.empty())
4264 {
4265 // Don't need hash entry.
4266 subMPT_.erase(simIterator);
4267 }
4268}
4269
4270void
4272{
4273 registry_.get().getJobQueue().addJob(JtClientAcctHist, "HistTxStream", [this, subInfo]() {
4274 auto const& accountId = subInfo.index->accountId;
4275 auto& lastLedgerSeq = subInfo.index->historyLastLedgerSeq;
4276 auto& txHistoryIndex = subInfo.index->historyTxIndex;
4277
4278 JLOG(journal_.trace()) << "AccountHistory job for account " << toBase58(accountId)
4279 << " started. lastLedgerSeq=" << lastLedgerSeq;
4280
4281 auto isFirstTx = [&](std::shared_ptr<Transaction> const& tx,
4282 std::shared_ptr<TxMeta> const& meta) -> bool {
4283 /*
4284 * genesis account: first tx is the one with seq 1
4285 * other account: first tx is the one created the account
4286 */
4287 if (accountId == kGenesisAccountId)
4288 {
4289 auto stx = tx->getSTransaction();
4290 if (stx->getAccountID(sfAccount) == accountId && stx->getSeqProxy().value() == 1)
4291 return true;
4292 }
4293
4294 return std::ranges::any_of(meta->getNodes(), [&](auto& node) {
4295 if (node.getFieldU16(sfLedgerEntryType) != ltACCOUNT_ROOT)
4296 return false;
4297
4298 if (node.isFieldPresent(sfNewFields))
4299 {
4300 if (auto inner = dynamic_cast<STObject const*>(node.peekAtPField(sfNewFields));
4301 inner)
4302 {
4303 if (inner->isFieldPresent(sfAccount) &&
4304 inner->getAccountID(sfAccount) == accountId)
4305 {
4306 return true;
4307 }
4308 }
4309 }
4310 return false;
4311 });
4312 };
4313
4314 auto send = [&](json::Value const& jvObj, bool unsubscribe) -> bool {
4315 if (auto sptr = subInfo.sinkWptr.lock())
4316 {
4317 sptr->send(jvObj, true);
4318 if (unsubscribe)
4319 unsubAccountHistory(sptr, accountId, false);
4320 return true;
4321 }
4322
4323 return false;
4324 };
4325
4326 auto sendMultiApiJson = [&](MultiApiJson const& jvObj, bool unsubscribe) -> bool {
4327 if (auto sptr = subInfo.sinkWptr.lock())
4328 {
4329 jvObj.visit(
4330 sptr->getApiVersion(), //
4331 [&](json::Value const& jv) { sptr->send(jv, true); });
4332
4333 if (unsubscribe)
4334 unsubAccountHistory(sptr, accountId, false);
4335 return true;
4336 }
4337
4338 return false;
4339 };
4340
4341 auto getMoreTxns = [&](std::uint32_t minLedger,
4342 std::uint32_t maxLedger,
4344 -> std::pair<
4347 auto& db = registry_.get().getRelationalDatabase();
4349 .account = accountId,
4350 .ledgerRange = {.min = minLedger, .max = maxLedger},
4351 .marker = marker,
4352 .limit = 0,
4353 .bAdmin = true,
4354 .delegate = std::nullopt};
4355 return db.newestAccountTxPage(options);
4356 };
4357
4358 /*
4359 * search backward until the genesis ledger or asked to stop
4360 */
4361 while (lastLedgerSeq >= 2 && !subInfo.index->stopHistorical)
4362 {
4363 int feeChargeCount = 0;
4364 if (auto sptr = subInfo.sinkWptr.lock(); sptr)
4365 {
4366 sptr->getConsumer().charge(resource::kFeeMediumBurdenRpc);
4367 ++feeChargeCount;
4368 }
4369 else
4370 {
4371 JLOG(journal_.trace())
4372 << "AccountHistory job for account " << toBase58(accountId)
4373 << " no InfoSub. Fee charged " << feeChargeCount << " times.";
4374 return;
4375 }
4376
4377 // try to search in 1024 ledgers till reaching genesis ledgers
4378 auto startLedgerSeq = (lastLedgerSeq > 1024 + 2 ? lastLedgerSeq - 1024 : 2);
4379 JLOG(journal_.trace())
4380 << "AccountHistory job for account " << toBase58(accountId)
4381 << ", working on ledger range [" << startLedgerSeq << "," << lastLedgerSeq << "]";
4382
4383 auto haveRange = [&]() -> bool {
4384 std::uint32_t validatedMin = UINT_MAX;
4385 std::uint32_t validatedMax = 0;
4386 auto haveSomeValidatedLedgers =
4387 registry_.get().getLedgerMaster().getValidatedRange(validatedMin, validatedMax);
4388
4389 return haveSomeValidatedLedgers && validatedMin <= startLedgerSeq &&
4390 lastLedgerSeq <= validatedMax;
4391 }();
4392
4393 if (!haveRange)
4394 {
4395 JLOG(journal_.debug()) << "AccountHistory reschedule job for account "
4396 << toBase58(accountId) << ", incomplete ledger range ["
4397 << startLedgerSeq << "," << lastLedgerSeq << "]";
4399 return;
4400 }
4401
4403 while (!subInfo.index->stopHistorical)
4404 {
4405 auto dbResult = getMoreTxns(startLedgerSeq, lastLedgerSeq, marker);
4406
4407 auto const& txns = dbResult.first;
4408 marker = dbResult.second;
4409 size_t const numTxns = txns.size();
4410 for (size_t i = 0; i < numTxns; ++i)
4411 {
4412 auto const& [tx, meta] = txns[i];
4413
4414 if (!tx || !meta)
4415 {
4416 JLOG(journal_.debug()) << "AccountHistory job for account "
4417 << toBase58(accountId) << " empty tx or meta.";
4418 send(rpcError(RpcInternal), true);
4419 return;
4420 }
4421 auto curTxLedger =
4422 registry_.get().getLedgerMaster().getLedgerBySeq(tx->getLedger());
4423 if (!curTxLedger)
4424 {
4425 // LCOV_EXCL_START
4426 UNREACHABLE(
4427 "xrpl::NetworkOPsImp::addAccountHistoryJob : "
4428 "getLedgerBySeq failed");
4429 JLOG(journal_.debug()) << "AccountHistory job for account "
4430 << toBase58(accountId) << " no ledger.";
4431 send(rpcError(RpcInternal), true);
4432 return;
4433 // LCOV_EXCL_STOP
4434 }
4435 std::shared_ptr<STTx const> const stTxn = tx->getSTransaction();
4436 if (!stTxn)
4437 {
4438 // LCOV_EXCL_START
4439 UNREACHABLE(
4440 "NetworkOPsImp::addAccountHistoryJob : "
4441 "getSTransaction failed");
4442 JLOG(journal_.debug()) << "AccountHistory job for account "
4443 << toBase58(accountId) << " getSTransaction failed.";
4444 send(rpcError(RpcInternal), true);
4445 return;
4446 // LCOV_EXCL_STOP
4447 }
4448
4449 auto const ref = std::ref(*meta);
4450 auto const trR = meta->getResultTER();
4451 MultiApiJson jvTx = transJson(stTxn, trR, true, curTxLedger, ref);
4452
4453 jvTx.set(jss::account_history_tx_index, txHistoryIndex--);
4454 if (i + 1 == numTxns || txns[i + 1].first->getLedger() != tx->getLedger())
4455 jvTx.set(jss::account_history_boundary, true);
4456
4457 if (isFirstTx(tx, meta))
4458 {
4459 jvTx.set(jss::account_history_tx_first, true);
4460 sendMultiApiJson(jvTx, false);
4461
4462 JLOG(journal_.trace()) << "AccountHistory job for account "
4463 << toBase58(accountId) << " done, found last tx.";
4464 return;
4465 }
4466
4467 sendMultiApiJson(jvTx, false);
4468 }
4469
4470 if (marker)
4471 {
4472 JLOG(journal_.trace())
4473 << "AccountHistory job for account " << toBase58(accountId)
4474 << " paging, marker=" << marker->ledgerSeq << ":" << marker->txnSeq;
4475 }
4476 else
4477 {
4478 break;
4479 }
4480 }
4481
4482 if (!subInfo.index->stopHistorical)
4483 {
4484 lastLedgerSeq = startLedgerSeq - 1;
4485 if (lastLedgerSeq <= 1)
4486 {
4487 JLOG(journal_.trace())
4488 << "AccountHistory job for account " << toBase58(accountId)
4489 << " done, reached genesis ledger.";
4490 return;
4491 }
4492 }
4493 }
4494 });
4495}
4496
4497void
4499 std::shared_ptr<ReadView const> const& ledger,
4501{
4502 subInfo.index->separationLedgerSeq = ledger->seq();
4503 auto const& accountId = subInfo.index->accountId;
4504 auto const accountKeylet = keylet::account(accountId);
4505 if (!ledger->exists(accountKeylet))
4506 {
4507 JLOG(journal_.debug()) << "subAccountHistoryStart, no account " << toBase58(accountId)
4508 << ", no need to add AccountHistory job.";
4509 return;
4510 }
4511 if (accountId == kGenesisAccountId)
4512 {
4513 if (auto const sleAcct = ledger->read(accountKeylet); sleAcct)
4514 {
4515 if (sleAcct->getFieldU32(sfSequence) == 1)
4516 {
4517 JLOG(journal_.debug())
4518 << "subAccountHistoryStart, genesis account " << toBase58(accountId)
4519 << " does not have tx, no need to add AccountHistory job.";
4520 return;
4521 }
4522 }
4523 else
4524 {
4525 // LCOV_EXCL_START
4526 UNREACHABLE(
4527 "xrpl::NetworkOPsImp::subAccountHistoryStart : failed to "
4528 "access genesis account");
4529 return;
4530 // LCOV_EXCL_STOP
4531 }
4532 }
4533 subInfo.index->historyLastLedgerSeq = ledger->seq();
4534 subInfo.index->haveHistorical = true;
4535
4536 JLOG(journal_.debug()) << "subAccountHistoryStart, add AccountHistory job: accountId="
4537 << toBase58(accountId) << ", currentLedgerSeq=" << ledger->seq();
4538
4539 addAccountHistoryJob(subInfo);
4540}
4541
4544{
4545 if (!isrListener->insertSubAccountHistory(accountId))
4546 {
4547 JLOG(journal_.debug()) << "subAccountHistory, already subscribed to account "
4548 << toBase58(accountId);
4549 return RpcInvalidParams;
4550 }
4551
4554 .sinkWptr = isrListener, .index = std::make_shared<SubAccountHistoryIndex>(accountId)};
4555 auto simIterator = subAccountHistory_.find(accountId);
4556 if (simIterator == subAccountHistory_.end())
4557 {
4559 inner.emplace(isrListener->getSeq(), ahi);
4560 subAccountHistory_.insert(simIterator, std::make_pair(accountId, inner));
4561 }
4562 else
4563 {
4564 simIterator->second.emplace(isrListener->getSeq(), ahi);
4565 }
4566
4567 auto const ledger = registry_.get().getLedgerMaster().getValidatedLedger();
4568 if (ledger)
4569 {
4570 subAccountHistoryStart(ledger, ahi);
4571 }
4572 else
4573 {
4574 // The node does not have validated ledgers, so wait for
4575 // one before start streaming.
4576 // In this case, the subscription is also considered successful.
4577 JLOG(journal_.debug()) << "subAccountHistory, no validated ledger yet, delay start";
4578 }
4579
4580 return RpcSuccess;
4581}
4582
4583void
4585 InfoSub::Ref isrListener,
4586 AccountID const& account,
4587 bool historyOnly)
4588{
4589 if (!historyOnly)
4590 isrListener->deleteSubAccountHistory(account);
4591 unsubAccountHistoryInternal(isrListener->getSeq(), account, historyOnly);
4592}
4593
4594void
4596 std::uint64_t seq,
4597 AccountID const& account,
4598 bool historyOnly)
4599{
4601 auto simIterator = subAccountHistory_.find(account);
4602 if (simIterator != subAccountHistory_.end())
4603 {
4604 auto& subInfoMap = simIterator->second;
4605 auto subInfoIter = subInfoMap.find(seq);
4606 if (subInfoIter != subInfoMap.end())
4607 {
4608 subInfoIter->second.index->stopHistorical = true;
4609 }
4610
4611 if (!historyOnly)
4612 {
4613 simIterator->second.erase(seq);
4614 if (simIterator->second.empty())
4615 {
4616 subAccountHistory_.erase(simIterator);
4617 }
4618 }
4619 JLOG(journal_.debug()) << "unsubAccountHistory, account " << toBase58(account)
4620 << ", historyOnly = " << (historyOnly ? "true" : "false");
4621 }
4622}
4623
4624bool
4626{
4627 // Server-side insert first, then InfoSub bookkeeping. If the InfoSub-side
4628 // insert throws, the orphan in subBook_ is cleared by the expired-weak_ptr
4629 // prune in pubBookTransaction. With the reverse ordering, ~InfoSub would
4630 // call unsubBookInternal for a key that was never inserted server-side.
4631 {
4632 std::scoped_lock const sl(bookLock_);
4633 subBook_[book].try_emplace(isrListener->getSeq(), isrListener);
4634 }
4635 isrListener->insertBookSubscription(book);
4636 return true;
4637}
4638
4639bool
4641{
4642 // Mirrors unsubAccount: clear the per-subscriber tracking set first so
4643 // ~InfoSub does not re-issue an unsubBookInternal for a book the caller
4644 // already removed, then erase the server-side entry.
4645 isrListener->deleteBookSubscription(book);
4646 return unsubBookInternal(isrListener->getSeq(), book);
4647}
4648
4649bool
4651{
4652 std::scoped_lock const sl(bookLock_);
4653 auto it = subBook_.find(book);
4654 if (it == subBook_.end())
4655 return false;
4656 bool const erased = it->second.erase(uSeq) != 0u;
4657 if (it->second.empty())
4658 subBook_.erase(it);
4659 return erased;
4660}
4661
4664{
4665 // This code-path is exclusively used when the server is in standalone
4666 // mode via `ledger_accept`
4667 XRPL_ASSERT(standalone_, "xrpl::NetworkOPsImp::acceptLedger : is standalone");
4668
4669 if (!standalone_)
4670 Throw<std::runtime_error>("Operation only possible in STANDALONE mode.");
4671
4672 // FIXME Could we improve on this and remove the need for a specialized
4673 // API in Consensus?
4674 beginConsensus(ledgerMaster_.getClosedLedger()->header().hash, {});
4675 consensus_.simulate(registry_.get().getTimeKeeper().closeTime(), consensusDelay);
4676 return ledgerMaster_.getCurrentLedger()->header().seq;
4677}
4678
4679// <-- bool: true=added, false=already there
4680bool
4682{
4683 if (auto lpClosed = ledgerMaster_.getValidatedLedger())
4684 {
4685 jvResult[jss::ledger_index] = lpClosed->header().seq;
4686 jvResult[jss::ledger_hash] = to_string(lpClosed->header().hash);
4687 jvResult[jss::ledger_time] =
4688 json::Value::UInt(lpClosed->header().closeTime.time_since_epoch().count());
4689 if (!lpClosed->rules().enabled(featureXRPFees))
4690 jvResult[jss::fee_ref] = kFeeUnitsDeprecated;
4691 jvResult[jss::fee_base] = lpClosed->fees().base.jsonClipped();
4692 jvResult[jss::reserve_base] = lpClosed->fees().reserve.jsonClipped();
4693 jvResult[jss::reserve_inc] = lpClosed->fees().increment.jsonClipped();
4694 jvResult[jss::network_id] = registry_.get().getNetworkIDService().getNetworkID();
4695 }
4696
4698 {
4699 jvResult[jss::validated_ledgers] = registry_.get().getLedgerMaster().getCompleteLedgers();
4700 }
4701
4702 std::scoped_lock const sl(streamLock_);
4703 return streamMaps_[SLedger].emplace(isrListener->getSeq(), isrListener).second;
4704}
4705
4706// <-- bool: true=added, false=already there
4707bool
4709{
4710 std::scoped_lock const sl(streamLock_);
4711 return streamMaps_[SBookChanges].emplace(isrListener->getSeq(), isrListener).second;
4712}
4713
4714// <-- bool: true=erased, false=was not there
4715bool
4717{
4718 std::scoped_lock const sl(streamLock_);
4719 return streamMaps_[SLedger].erase(uSeq) != 0u;
4720}
4721
4722// <-- bool: true=erased, false=was not there
4723bool
4725{
4726 std::scoped_lock const sl(streamLock_);
4727 return streamMaps_[SBookChanges].erase(uSeq) != 0u;
4728}
4729
4730// <-- bool: true=added, false=already there
4731bool
4733{
4734 std::scoped_lock const sl(streamLock_);
4735 return streamMaps_[SManifests].emplace(isrListener->getSeq(), isrListener).second;
4736}
4737
4738// <-- bool: true=erased, false=was not there
4739bool
4741{
4742 std::scoped_lock const sl(streamLock_);
4743 return streamMaps_[SManifests].erase(uSeq) != 0u;
4744}
4745
4746// <-- bool: true=added, false=already there
4747bool
4748NetworkOPsImp::subServer(InfoSub::Ref isrListener, json::Value& jvResult, bool admin)
4749{
4750 UInt256 uRandom;
4751
4752 if (standalone_)
4753 jvResult[jss::stand_alone] = standalone_;
4754
4755 // CHECKME: is it necessary to provide a random number here?
4756 beast::rngfill(uRandom.begin(), uRandom.size(), cryptoPrng());
4757
4758 auto const& feeTrack = registry_.get().getFeeTrack();
4759 jvResult[jss::random] = to_string(uRandom);
4760 jvResult[jss::server_status] = strOperatingMode(admin);
4761 jvResult[jss::load_base] = feeTrack.getLoadBase();
4762 jvResult[jss::load_factor] = feeTrack.getLoadFactor();
4763 jvResult[jss::hostid] = getHostId(admin);
4764 jvResult[jss::pubkey_node] =
4765 toBase58(TokenType::NodePublic, registry_.get().getApp().nodeIdentity().first);
4766
4767 std::scoped_lock const sl(streamLock_);
4768 return streamMaps_[SServer].emplace(isrListener->getSeq(), isrListener).second;
4769}
4770
4771// <-- bool: true=erased, false=was not there
4772bool
4774{
4775 std::scoped_lock const sl(streamLock_);
4776 return streamMaps_[SServer].erase(uSeq) != 0u;
4777}
4778
4779// <-- bool: true=added, false=already there
4780bool
4782{
4783 std::scoped_lock const sl(streamLock_);
4784 return streamMaps_[STransactions].emplace(isrListener->getSeq(), isrListener).second;
4785}
4786
4787// <-- bool: true=erased, false=was not there
4788bool
4790{
4791 std::scoped_lock const sl(streamLock_);
4792 return streamMaps_[STransactions].erase(uSeq) != 0u;
4793}
4794
4795// <-- bool: true=added, false=already there
4796bool
4798{
4799 std::scoped_lock const sl(streamLock_);
4800 return streamMaps_[SRtTransactions].emplace(isrListener->getSeq(), isrListener).second;
4801}
4802
4803// <-- bool: true=erased, false=was not there
4804bool
4806{
4807 std::scoped_lock const sl(streamLock_);
4808 return streamMaps_[SRtTransactions].erase(uSeq) != 0u;
4809}
4810
4811// <-- bool: true=added, false=already there
4812bool
4814{
4815 std::scoped_lock const sl(streamLock_);
4816 return streamMaps_[SValidations].emplace(isrListener->getSeq(), isrListener).second;
4817}
4818
4819void
4824
4825// <-- bool: true=erased, false=was not there
4826bool
4828{
4829 std::scoped_lock const sl(streamLock_);
4830 return streamMaps_[SValidations].erase(uSeq) != 0u;
4831}
4832
4833// <-- bool: true=added, false=already there
4834bool
4836{
4837 std::scoped_lock const sl(streamLock_);
4838 return streamMaps_[SPeerStatus].emplace(isrListener->getSeq(), isrListener).second;
4839}
4840
4841// <-- bool: true=erased, false=was not there
4842bool
4844{
4845 std::scoped_lock const sl(streamLock_);
4846 return streamMaps_[SPeerStatus].erase(uSeq) != 0u;
4847}
4848
4849// <-- bool: true=added, false=already there
4850bool
4852{
4853 std::scoped_lock const sl(streamLock_);
4854 return streamMaps_[SConsensusPhase].emplace(isrListener->getSeq(), isrListener).second;
4855}
4856
4857// <-- bool: true=erased, false=was not there
4858bool
4860{
4861 std::scoped_lock const sl(streamLock_);
4862 return streamMaps_[SConsensusPhase].erase(uSeq) != 0u;
4863}
4864
4867{
4868 // Caller already holds streamLock_; this performs the lookup only.
4869 auto const it = rpcSubMap_.find(strUrl);
4870
4871 if (it != rpcSubMap_.end())
4872 return it->second;
4873
4874 return InfoSub::pointer();
4875}
4876
4879{
4880 std::scoped_lock const sl(streamLock_);
4881 return findRpcSubLocked(strUrl);
4882}
4883
4886{
4887 std::scoped_lock const sl(streamLock_);
4888
4889 rpcSubMap_.emplace(strUrl, rspEntry);
4890
4891 return rspEntry;
4892}
4893
4894bool
4896{
4897 // Declared before the lock so it outlives the scoped_lock and is destroyed
4898 // only after streamLock_ is released. The erase below may drop the last
4899 // strong reference; if so, ~InfoSub runs and its unsub* calls re-acquire
4900 // the non-recursive streamLock_. Destroying pInfo inside the lock would
4901 // self-deadlock.
4902 InfoSub::pointer pInfo;
4903 {
4904 std::scoped_lock const sl(streamLock_);
4905 // Use the no-lock helper: we already hold streamLock_ and the mutex is
4906 // not recursive, so calling the public findRpcSub here would deadlock.
4907 pInfo = findRpcSubLocked(strUrl);
4908
4909 if (!pInfo)
4910 return false;
4911
4912 // check to see if any of the stream maps still hold a weak reference to
4913 // this entry before removing
4914 for (SubMapType const& map : streamMaps_)
4915 {
4916 if (map.contains(pInfo->getSeq()))
4917 return false;
4918 }
4919 rpcSubMap_.erase(strUrl);
4920 }
4921 // pInfo destroyed here, after streamLock_ is released.
4922 return true;
4923}
4924
4925#ifndef USE_NEW_BOOK_PAGE
4926
4927// NIKB FIXME this should be looked at. There's no reason why this shouldn't
4928// work, but it demonstrated poor performance.
4929//
4930void
4933 Book const& book,
4934 AccountID const& uTakerID,
4935 bool const bProof,
4936 unsigned int iLimit,
4937 json::Value const& jvMarker,
4938 json::Value& jvResult)
4939{ // CAUTION: This is the old get book page logic
4940 json::Value& jvOffers = (jvResult[jss::offers] = json::Value(json::ValueType::Array));
4941
4943 UInt256 const uBookBase = getBookBase(book);
4944 UInt256 const uBookEnd = getQualityNext(uBookBase);
4945 UInt256 uTipIndex = uBookBase;
4946
4947 if (auto stream = journal_.trace())
4948 {
4949 stream << "getBookPage:" << book;
4950 stream << "getBookPage: uBookBase=" << uBookBase;
4951 stream << "getBookPage: uBookEnd=" << uBookEnd;
4952 stream << "getBookPage: uTipIndex=" << uTipIndex;
4953 }
4954
4955 ReadView const& view = *lpLedger;
4956
4957 bool const bGlobalFreeze = isGlobalFrozen(view, book.out) || isGlobalFrozen(view, book.in);
4958
4959 bool bDone = false;
4960 bool bDirectAdvance = true;
4961
4962 SLE::const_pointer sleOfferDir;
4963 UInt256 offerIndex;
4964 unsigned int uBookEntry = 0;
4965 STAmount saDirRate;
4966
4967 auto const rate = transferRate(view, book.out);
4968 auto viewJ = registry_.get().getJournal("View");
4969
4970 while (!bDone && iLimit-- > 0)
4971 {
4972 if (bDirectAdvance)
4973 {
4974 bDirectAdvance = false;
4975
4976 JLOG(journal_.trace()) << "getBookPage: bDirectAdvance";
4977
4978 auto const ledgerIndex = view.succ(uTipIndex, uBookEnd);
4979 if (ledgerIndex)
4980 {
4981 sleOfferDir = view.read(keylet::page(*ledgerIndex));
4982 }
4983 else
4984 {
4985 sleOfferDir.reset();
4986 }
4987
4988 if (!sleOfferDir)
4989 {
4990 JLOG(journal_.trace()) << "getBookPage: bDone";
4991 bDone = true;
4992 }
4993 else
4994 {
4995 uTipIndex = sleOfferDir->key();
4996 saDirRate = amountFromQuality(getQuality(uTipIndex));
4997
4998 cdirFirst(view, uTipIndex, sleOfferDir, uBookEntry, offerIndex);
4999
5000 JLOG(journal_.trace()) << "getBookPage: uTipIndex=" << uTipIndex;
5001 JLOG(journal_.trace()) << "getBookPage: offerIndex=" << offerIndex;
5002 }
5003 }
5004
5005 if (!bDone)
5006 {
5007 auto sleOffer = view.read(keylet::offer(offerIndex));
5008
5009 if (sleOffer)
5010 {
5011 auto const uOfferOwnerID = sleOffer->getAccountID(sfAccount);
5012 auto const& saTakerGets = sleOffer->getFieldAmount(sfTakerGets);
5013 auto const& saTakerPays = sleOffer->getFieldAmount(sfTakerPays);
5014 STAmount saOwnerFunds;
5015 bool firstOwnerOffer(true);
5016 auto foundBalance = [&]() {
5017 auto umBalanceEntry = umBalance.find(uOfferOwnerID);
5018 if (umBalanceEntry == umBalance.end())
5019 return false;
5020
5021 // Found in running balance table.
5022 saOwnerFunds = umBalanceEntry->second;
5023 firstOwnerOffer = false;
5024 return true;
5025 };
5026
5027 if (book.out.getIssuer() == uOfferOwnerID)
5028 {
5029 book.out.visit(
5030 [&](Issue const&) {
5031 // If an offer is selling issuer's own IOUs, it is
5032 // fully funded.
5033 saOwnerFunds = saTakerGets;
5034 },
5035 [&](MPTIssue const& issue) {
5036 // MPT issuers have bounded self-issuance. Use the
5037 // running balance table so multiple issuer-owned
5038 // offers share the same remaining issuance
5039 // headroom.
5040 if (!foundBalance())
5041 {
5042 // Did not find balance in table.
5043
5044 saOwnerFunds = issuerFundsToSelfIssue(view, issue);
5045 }
5046 });
5047 }
5048 else if (bGlobalFreeze)
5049 {
5050 // If either asset is globally frozen, consider all offers
5051 // that aren't ours to be totally unfunded
5052 saOwnerFunds.clear(book.out);
5053 }
5054 else
5055 {
5056 if (!foundBalance())
5057 {
5058 // Did not find balance in table.
5059
5060 saOwnerFunds = accountHolds(
5061 view,
5062 uOfferOwnerID,
5063 book.out,
5066 viewJ);
5067
5068 if (saOwnerFunds < beast::kZero)
5069 {
5070 // Treat negative funds as zero.
5071
5072 saOwnerFunds.clear();
5073 }
5074 }
5075 }
5076
5077 json::Value jvOffer = sleOffer->getJson(JsonOptions::Values::None);
5078
5079 STAmount saTakerGetsFunded;
5080 STAmount saOwnerFundsLimit = saOwnerFunds;
5081 Rate offerRate = kParityRate;
5082
5083 if (rate != kParityRate
5084 // Have a transfer fee.
5085 && uTakerID != book.out.getIssuer()
5086 // Not taking offers of own IOUs.
5087 && book.out.getIssuer() != uOfferOwnerID)
5088 // Offer owner not issuing ownfunds
5089 {
5090 // Need to charge a transfer fee to offer owner.
5091 offerRate = rate;
5092 // Why MPT does not use divide(): divide() is built for an
5093 // IOU mantissa, which is always normalized into
5094 // [1e15, 1e16]. An MPT mantissa is the raw int64 balance,
5095 // and divide() scales the numerator by 1e17, so a balance
5096 // over ~1.8e17 leaves uint64 range and throws -- failing
5097 // the whole RPC rather than this one offer.
5098 //
5099 // Why mulRatio is safe: it evaluates in 128 bits, and here
5100 // it cannot overflow either. offerRate is
5101 // 1e9 + 10'000 * TransferFee, so this branch runs only with
5102 // offerRate > kParityRate, making the quotient smaller than
5103 // saOwnerFunds. Rounded down, so reported liquidity is
5104 // never overstated.
5105 saOwnerFundsLimit = saOwnerFunds.holds<MPTIssue>()
5106 ? toSTAmount(
5107 mulRatio(
5108 saOwnerFunds.mpt(),
5109 kParityRate.value,
5110 offerRate.value,
5111 /*roundUp*/ false),
5112 saOwnerFunds.asset())
5113 : divide(saOwnerFunds, offerRate);
5114 }
5115
5116 if (saOwnerFundsLimit >= saTakerGets)
5117 {
5118 // Sufficient funds no shenanigans.
5119 saTakerGetsFunded = saTakerGets;
5120 }
5121 else
5122 {
5123 // Only provide, if not fully funded.
5124
5125 saTakerGetsFunded = saOwnerFundsLimit;
5126
5127 saTakerGetsFunded.setJson(jvOffer[jss::taker_gets_funded]);
5128 std::min(
5129 saTakerPays, multiply(saTakerGetsFunded, saDirRate, saTakerPays.asset()))
5130 .setJson(jvOffer[jss::taker_pays_funded]);
5131 }
5132
5133 // What the owner pays for this offer, charged against the
5134 // running balance for the owner's later offers. BookStep
5135 // rounds the MPT owner's fee-grossed cost up; multiply() rounds
5136 // to nearest and would leave one unit per offer over-reported
5137 // to the next offer. Can't overflow: saTakerGetsFunded <=
5138 // floor(saOwnerFunds / offerRate).
5139 auto const grossed = [&]() {
5140 return saOwnerFunds.asset().visit(
5141 [&](MPTIssue const&) {
5142 auto const mpt = mulRatio(
5143 saTakerGetsFunded.mpt(),
5144 offerRate.value,
5145 kParityRate.value,
5146 /*roundUp*/ true);
5147 return toSTAmount(mpt, saOwnerFunds.asset());
5148 },
5149 [&](Issue const&) { return multiply(saTakerGetsFunded, offerRate); });
5150 };
5151 STAmount const saOwnerPays = (kParityRate == offerRate)
5152 ? saTakerGetsFunded
5153 : std::min(saOwnerFunds, grossed());
5154
5155 umBalance[uOfferOwnerID] = saOwnerFunds - saOwnerPays;
5156
5157 // Include all offers funded and unfunded
5158 json::Value& jvOf = jvOffers.append(jvOffer);
5159 jvOf[jss::quality] = saDirRate.getText();
5160
5161 if (firstOwnerOffer)
5162 jvOf[jss::owner_funds] = saOwnerFunds.getText();
5163 }
5164 else
5165 {
5166 JLOG(journal_.warn()) << "Missing offer";
5167 }
5168
5169 if (!cdirNext(view, uTipIndex, sleOfferDir, uBookEntry, offerIndex))
5170 {
5171 bDirectAdvance = true;
5172 }
5173 else
5174 {
5175 JLOG(journal_.trace()) << "getBookPage: offerIndex=" << offerIndex;
5176 }
5177 }
5178 }
5179
5180 // jvResult[jss::marker] = json::Value(json::ValueType::Array);
5181 // jvResult[jss::nodes] = json::Value(json::ValueType::Array);
5182}
5183
5184#else
5185
5186// This is the new code that uses the book iterators
5187// It has temporarily been disabled
5188// If this path is re-enabled, add MPT support mirroring the functional
5189// getBookPage() implementation above.
5190
5191void
5194 Book const& book,
5195 AccountID const& uTakerID,
5196 bool const bProof,
5197 unsigned int iLimit,
5198 json::Value const& jvMarker,
5199 json::Value& jvResult)
5200{
5201 auto& jvOffers = (jvResult[jss::offers] = json::Value(json::ValueType::Array));
5202
5204
5205 MetaView lesActive(lpLedger, tapNONE, true);
5206 OrderBookIterator obIterator(lesActive, book);
5207
5208 auto const rate = transferRate(lesActive, book.out.account);
5209
5210 bool const bGlobalFreeze =
5211 lesActive.isGlobalFrozen(book.out.account) || lesActive.isGlobalFrozen(book.in.account);
5212
5213 while (iLimit-- > 0 && obIterator.nextOffer())
5214 {
5215 SLE::pointer sleOffer = obIterator.getCurrentOffer();
5216 if (sleOffer)
5217 {
5218 auto const uOfferOwnerID = sleOffer->getAccountID(sfAccount);
5219 auto const& saTakerGets = sleOffer->getFieldAmount(sfTakerGets);
5220 auto const& saTakerPays = sleOffer->getFieldAmount(sfTakerPays);
5221 STAmount saDirRate = obIterator.getCurrentRate();
5222 STAmount saOwnerFunds;
5223
5224 if (book.out.account == uOfferOwnerID)
5225 {
5226 // If offer is selling issuer's own IOUs, it is fully funded.
5227 saOwnerFunds = saTakerGets;
5228 }
5229 else if (bGlobalFreeze)
5230 {
5231 // If either asset is globally frozen, consider all offers
5232 // that aren't ours to be totally unfunded
5233 saOwnerFunds.clear(book.out);
5234 }
5235 else
5236 {
5237 auto umBalanceEntry = umBalance.find(uOfferOwnerID);
5238
5239 if (umBalanceEntry != umBalance.end())
5240 {
5241 // Found in running balance table.
5242
5243 saOwnerFunds = umBalanceEntry->second;
5244 }
5245 else
5246 {
5247 // Did not find balance in table.
5248
5249 saOwnerFunds = lesActive.accountHolds(
5250 uOfferOwnerID,
5251 book.out.currency,
5252 book.out.account,
5254
5255 if (saOwnerFunds.isNegative())
5256 {
5257 // Treat negative funds as zero.
5258
5259 saOwnerFunds.zero();
5260 }
5261 }
5262 }
5263
5264 json::Value jvOffer = sleOffer->getJson(JsonOptions::Values::None);
5265
5266 STAmount saTakerGetsFunded;
5267 STAmount saOwnerFundsLimit = saOwnerFunds;
5268 Rate offerRate = parityRate;
5269
5270 if (rate != parityRate
5271 // Have a transfer fee.
5272 && uTakerID != book.out.account
5273 // Not taking offers of own IOUs.
5274 && book.out.account != uOfferOwnerID)
5275 // Offer owner not issuing ownfunds
5276 {
5277 // Need to charge a transfer fee to offer owner.
5278 offerRate = rate;
5279 saOwnerFundsLimit = divide(saOwnerFunds, offerRate);
5280 }
5281
5282 if (saOwnerFundsLimit >= saTakerGets)
5283 {
5284 // Sufficient funds no shenanigans.
5285 saTakerGetsFunded = saTakerGets;
5286 }
5287 else
5288 {
5289 // Only provide, if not fully funded.
5290 saTakerGetsFunded = saOwnerFundsLimit;
5291
5292 saTakerGetsFunded.setJson(jvOffer[jss::taker_gets_funded]);
5293
5294 // TODO(tom): The result of this expression is not used - what's
5295 // going on here?
5296 std::min(saTakerPays, multiply(saTakerGetsFunded, saDirRate, saTakerPays.asset()))
5297 .setJson(jvOffer[jss::taker_pays_funded]);
5298 }
5299
5300 STAmount saOwnerPays = (parityRate == offerRate)
5301 ? saTakerGetsFunded
5302 : std::min(saOwnerFunds, multiply(saTakerGetsFunded, offerRate));
5303
5304 umBalance[uOfferOwnerID] = saOwnerFunds - saOwnerPays;
5305
5306 if (!saOwnerFunds.isZero() || uOfferOwnerID == uTakerID)
5307 {
5308 // Only provide funded offers and offers of the taker.
5309 json::Value& jvOf = jvOffers.append(jvOffer);
5310 jvOf[jss::quality] = saDirRate.getText();
5311 }
5312 }
5313 }
5314
5315 // jvResult[jss::marker] = json::Value(json::ValueType::Array);
5316 // jvResult[jss::nodes] = json::Value(json::ValueType::Array);
5317}
5318
5319#endif
5320
5321inline void
5323{
5324 auto [counters, mode, start, initialSync] = accounting_.getCounterData();
5327 counters[static_cast<std::size_t>(mode)].dur += current;
5328
5329 std::scoped_lock const lock(statsMutex_);
5330 stats_.disconnectedDuration.set(
5331 counters[static_cast<std::size_t>(OperatingMode::DISCONNECTED)].dur.count());
5332 stats_.connectedDuration.set(
5333 counters[static_cast<std::size_t>(OperatingMode::CONNECTED)].dur.count());
5334 stats_.syncingDuration.set(
5335 counters[static_cast<std::size_t>(OperatingMode::SYNCING)].dur.count());
5336 stats_.trackingDuration.set(
5337 counters[static_cast<std::size_t>(OperatingMode::TRACKING)].dur.count());
5338 stats_.fullDuration.set(counters[static_cast<std::size_t>(OperatingMode::FULL)].dur.count());
5339
5340 stats_.disconnectedTransitions.set(
5341 counters[static_cast<std::size_t>(OperatingMode::DISCONNECTED)].transitions);
5342 stats_.connectedTransitions.set(
5343 counters[static_cast<std::size_t>(OperatingMode::CONNECTED)].transitions);
5344 stats_.syncingTransitions.set(
5345 counters[static_cast<std::size_t>(OperatingMode::SYNCING)].transitions);
5346 stats_.trackingTransitions.set(
5347 counters[static_cast<std::size_t>(OperatingMode::TRACKING)].transitions);
5348 stats_.fullTransitions.set(counters[static_cast<std::size_t>(OperatingMode::FULL)].transitions);
5349}
5350
5351void
5353{
5354 auto now = std::chrono::steady_clock::now();
5355
5356 std::scoped_lock const lock(mutex_);
5357 ++counters_[static_cast<std::size_t>(om)].transitions;
5358 if (om == OperatingMode::FULL && counters_[static_cast<std::size_t>(om)].transitions == 1)
5359 {
5362 }
5363 counters_[static_cast<std::size_t>(mode_)].dur +=
5365
5366 mode_ = om;
5367 start_ = now;
5368}
5369
5370void
5372{
5373 auto [counters, mode, start, initialSync] = getCounterData();
5376 counters[static_cast<std::size_t>(mode)].dur += current;
5377
5378 obj[jss::state_accounting] = json::ValueType::Object;
5379 for (auto i = static_cast<std::size_t>(OperatingMode::DISCONNECTED);
5380 i <= static_cast<std::size_t>(OperatingMode::FULL);
5381 ++i)
5382 {
5383 obj[jss::state_accounting][kStates[i]] = json::ValueType::Object;
5384 auto& state = obj[jss::state_accounting][kStates[i]];
5385 state[jss::transitions] = std::to_string(counters[i].transitions);
5386 state[jss::duration_us] = std::to_string(counters[i].dur.count());
5387 }
5388 obj[jss::server_state_duration_us] = std::to_string(current.count());
5389 if (initialSync != 0u)
5390 obj[jss::initial_sync_duration_us] = std::to_string(initialSync);
5391}
5392
5393//------------------------------------------------------------------------------
5394
5397 ServiceRegistry& registry,
5398 NetworkOPs::ClockType& clock,
5399 bool standalone,
5400 std::size_t minPeerCount,
5401 bool startValid,
5402 JobQueue& jobQueue,
5403 LedgerMaster& ledgerMaster,
5404 ValidatorKeys const& validatorKeys,
5405 boost::asio::io_context& ioCtx,
5407 beast::insight::Collector::Ptr const& collector)
5408{
5410 registry,
5411 clock,
5412 standalone,
5413 minPeerCount,
5414 startValid,
5415 jobQueue,
5416 ledgerMaster,
5417 validatorKeys,
5418 ioCtx,
5419 journal,
5420 collector);
5421}
5422
5423} // namespace xrpl
T any_of(T... args)
T back_inserter(T... args)
T begin(T... args)
A generic endpoint for log messages.
Definition Journal.h:44
std::shared_ptr< Collector > Ptr
Definition Collector.h:29
A metric for measuring an integral value.
Definition Gauge.h:21
A reference to a handler for performing polled collection.
Definition Hook.h:14
Decorator for streaming out compact json.
Lightweight wrapper to tag static string.
Definition json_value.h:48
Represents a JSON value.
Definition json_value.h:117
Value get(UInt index, Value const &defaultValue) const
If the array contains at least index+1 elements, returns the element value, otherwise returns default...
json::UInt UInt
Definition json_value.h:124
Value & append(Value const &value)
Append value to array at the end.
UInt size() const
Number of values in array or object.
bool isMember(char const *key) const
Return true if the object has a member named key.
A transaction that is in a closed ledger.
TxMeta const & getMeta() const
boost::container::flat_set< AccountID > const & getAffected() const
std::shared_ptr< STTx const > const & getTxn() const
constexpr auto visit(Visitors &&... visitors) const -> decltype(auto)
Definition Asset.h:117
bool isZero() const
Definition base_uint.h:562
bool isNonZero() const
Definition base_uint.h:567
iterator begin()
Definition base_uint.h:128
static constexpr std::size_t size()
Definition base_uint.h:548
Specifies an order book.
Definition Book.h:28
Holds transactions which were deferred to the next pass of consensus.
The role of a ClosureCounter is to assist in shutdown by letting callers wait for the completion of c...
std::uint32_t getLoadFee() const
Definition ClusterNode.h:33
NetClock::time_point getReportTime() const
Definition ClusterNode.h:39
PublicKey const & identity() const
Definition ClusterNode.h:45
std::string const & name() const
Definition ClusterNode.h:27
std::shared_ptr< InfoSub > const & Ref
Definition InfoSub.h:98
std::shared_ptr< InfoSub > pointer
Definition InfoSub.h:92
std::weak_ptr< InfoSub > Wptr
Definition InfoSub.h:96
A currency issued by an account.
Definition Issue.h:18
A pool of threads to perform work.
Definition JobQueue.h:60
Tracks the current ledger and any ledgers in the process of closing.
Manages the current fee schedule.
Manages load sources.
Definition LoadManager.h:28
void heartbeat()
Reset the stall detection timer.
static constexpr int kHoldLedgers
Definition LocalTxs.h:23
State accounting records two attributes for each possible server state: 1) Amount of time spent in ea...
void json(json::Value &obj) const
Output state counters in JSON format.
std::chrono::steady_clock::time_point const processStart_
static std::array< json::StaticString const, 5 > const kStates
std::array< Counters, 5 > counters_
void mode(OperatingMode om)
Record state transition.
std::chrono::steady_clock::time_point start_
Transaction with input flags and results to be applied in batches.
std::shared_ptr< Transaction > const transaction
TransactionStatus(std::shared_ptr< Transaction > t, bool a, bool l, FailHard f)
SubBookMapType subBook_
Guarded by bookLock_.
std::string getHostId(bool forAdmin)
void reportConsensusStateChange(ConsensusPhase phase)
void addAccountHistoryJob(SubAccountHistoryInfoWeak subInfo)
void clearNeedNetworkLedger() override
SubInfoMapType subAccount_
HashMap< AccountID, HashMap< std::uint64_t, SubAccountHistoryInfoWeak > > SubAccountHistoryMapType
static constexpr std::size_t kAccountCleanupChunk
Maximum number of account entries erased per accountLock_ acquisition during disconnect-time cleanup.
bool subManifests(InfoSub::Ref ispListener) override
std::size_t const minPeerCount_
bool subServer(InfoSub::Ref ispListener, json::Value &jvResult, bool admin) override
std::vector< TransactionStatus > transactions_
bool subPeerStatus(InfoSub::Ref ispListener) override
void pubAccountTransaction(std::shared_ptr< ReadView const > const &ledger, AcceptedLedgerTx const &transaction, bool last)
void subMPT(InfoSub::Ref ispListener, HashSet< MPTID > const &mptIDs) override
void unsubAccount(InfoSub::Ref ispListener, HashSet< AccountID > const &vnaAccountIDs, bool rt) override
bool subValidations(InfoSub::Ref ispListener) override
ClosureCounter< void, boost::system::error_code const & > waitHandlerCounter_
std::condition_variable cond_
MultiApiJson transJson(std::shared_ptr< STTx const > const &transaction, TER result, bool validated, std::shared_ptr< ReadView const > const &ledger, std::optional< std::reference_wrapper< TxMeta const > > meta)
bool unsubManifests(std::uint64_t uListener) override
json::Value getOwnerInfo(std::shared_ptr< ReadView const > lpLedger, AccountID const &account) override
void unsubMPTInternal(std::uint64_t seq, MPTID const &mptID) override
Remove an MPT issuance subscription during InfoSub teardown.
HashMap< Book, SubMapType > SubBookMapType
Maps each order book to its current set of subscribers.
void cleanupAccountHistorySubscriptions(std::uint64_t seq, HashSet< AccountID > const &accounts)
Erase one connection's entries from subAccountHistory_ in accountLock_-bounded chunks.
void stateAccounting(json::Value &obj) override
void pubLedger(std::shared_ptr< ReadView const > const &lpAccepted) override
void stop() override
DispatchState dispatchState_
InfoSub::pointer addRpcSub(std::string const &strUrl, InfoSub::Ref) override
void transactionBatch()
Apply transactions in batches.
void setTimer(boost::asio::steady_timer &timer, std::chrono::milliseconds const &expiryTime, std::function< void()> onExpire, std::function< void()> onError)
bool unsubRTTransactions(std::uint64_t uListener) override
beast::Journal const & journal() const override
Journal used by InfoSub for diagnostics that occur after the owning subsystem (e.g.
json::Value getLedgerFetchInfo() override
bool processTrustedProposal(RCLCxPeerPos proposal) override
InfoSub::pointer findRpcSubLocked(std::string const &strUrl)
Look up an RPC subscription without taking streamLock_.
void kickoffAccountHistory(std::shared_ptr< AcceptedLedger const > const &alpAccepted)
On the first published ledger only, start the delayed account-history streaming for any subscriptions...
void subAccountHistoryStart(std::shared_ptr< ReadView const > const &ledger, SubAccountHistoryInfoWeak &subInfo)
void pubValidation(std::shared_ptr< STValidation > const &val) override
json::Value getConsensusInfo() override
ConsensusPhase lastConsensusPhase_
beast::Journal journal_
void endConsensus(std::unique_ptr< std::stringstream > const &clog) override
bool beginConsensus(UInt256 const &networkClosed, std::unique_ptr< std::stringstream > const &clog) override
void setMode(OperatingMode om) override
void setAmendmentBlocked() override
bool subBook(InfoSub::Ref ispListener, Book const &) override
ErrorCodeI subAccountHistory(InfoSub::Ref ispListener, AccountID const &account) override
subscribe an account's new transactions and retrieve the account's historical transactions
HashMap< AccountID, SubMapType > SubInfoMapType
void pubConsensus(ConsensusPhase phase)
bool isNeedNetworkLedger() override
bool unsubBookInternal(std::uint64_t uListener, Book const &) override
Remove a book subscription during InfoSub teardown.
DispatchState
Synchronization states for transaction batches.
SubInfoMapType subRTAccount_
std::atomic< bool > needNetworkLedger_
boost::asio::steady_timer heartbeatTimer_
std::mutex mptLock_
Guards subMPT_.
std::array< SubMapType, SubTypes::SLastEntry > streamMaps_
One weak_ptr subscriber map per stream type.
std::reference_wrapper< ServiceRegistry > registry_
static std::array< char const *, 5 > const kStates
bool unsubLedger(std::uint64_t uListener) override
void pubProposedAccountTransaction(std::shared_ptr< ReadView const > const &ledger, std::shared_ptr< STTx const > const &transaction, TER result)
void unsubAccountHistoryInternal(std::uint64_t seq, AccountID const &account, bool historyOnly) override
void pubValidatedTransaction(std::shared_ptr< ReadView const > const &ledger, AcceptedLedgerTx const &transaction, bool last)
void pubBookTransaction(AcceptedLedgerTx const &transaction, MultiApiJson const &jvObj)
Fan transaction notifications out to all book subscribers.
void cleanupAccountSubscriptions(std::uint64_t seq, HashSet< AccountID > const &accounts, SubInfoMapType &subMap)
Erase one connection's entries from the given account map (subAccount_ or subRTAccount_) in accountLo...
void switchLastClosedLedger(std::shared_ptr< Ledger const > const &newLCL)
std::optional< PublicKey > const validatorPK_
std::mutex streamLock_
Guards streamMaps_[] and rpcSubMap_.
std::atomic< bool > amendmentBlocked_
bool unsubBook(InfoSub::Ref ispListener, Book const &) override
Remove a book subscription for a live subscriber.
bool isFull() override
void clearAmendmentWarned() override
LedgerMaster & ledgerMaster_
void publishLedgerStreams(std::shared_ptr< ReadView const > const &lpAccepted, std::shared_ptr< AcceptedLedger const > const &alpAccepted)
Send the ledgerClosed and book-changes stream updates for a ledger.
std::atomic< OperatingMode > mode_
void updateLocalTx(ReadView const &view) override
void clearLedgerFetch() override
std::unique_ptr< LocalTxs > localTX_
void apply(std::unique_lock< std::mutex > &batchLock)
Attempt to apply transactions and post-process based on the results.
InfoSub::pointer findRpcSub(std::string const &strUrl) override
bool isAmendmentBlocked() override
std::string strOperatingMode(OperatingMode const mode, bool const admin) const override
bool subBookChanges(InfoSub::Ref ispListener) override
SubRpcMapType rpcSubMap_
HashMap< std::string, InfoSub::pointer > SubRpcMapType
void setStandAlone() override
void setNeedNetworkLedger() override
HashMap< MPTID, SubMapType > SubMPTInfoMapType
std::size_t getBookSubscribersCount() override
Total number of (book, subscriber) entries currently tracked.
std::mutex bookLock_
Guards subBook_.
bool unsubServer(std::uint64_t uListener) override
bool unsubConsensus(std::uint64_t uListener) override
bool subConsensus(InfoSub::Ref ispListener) override
void pubManifest(Manifest const &) override
std::mutex accountLock_
Guards subAccount_, subRTAccount_, subAccountHistory_.
void consensusViewChange() override
boost::asio::steady_timer accountHistoryTxTimer_
NetworkOPsImp(ServiceRegistry &registry, NetworkOPs::ClockType &clock, bool standalone, std::size_t minPeerCount, bool startValid, JobQueue &jobQueue, LedgerMaster &ledgerMaster, ValidatorKeys const &validatorKeys, boost::asio::io_context &ioCtx, beast::Journal journal, beast::insight::Collector::Ptr const &collector)
bool recvValidation(std::shared_ptr< STValidation > const &val, std::string const &source) override
void setUNLBlocked() override
bool unsubValidations(std::uint64_t uListener) override
void scheduleAccountCleanup(std::uint64_t seq, HashSet< AccountID > rtAccounts, HashSet< AccountID > normalAccounts, HashSet< AccountID > historyAccounts) override
Schedule the server-side teardown of a disconnecting connection's account subscriptions off the destr...
void pubMPTTransaction(AcceptedLedgerTx const &transaction, MultiApiJson const &jvObj)
void doTransactionAsync(std::shared_ptr< Transaction > transaction, bool bUnlimited, FailHard failtype)
For transactions not submitted by a locally connected client, fire and forget.
std::set< UInt256 > pendingValidations_
bool subRTTransactions(InfoSub::Ref ispListener) override
OperatingMode getOperatingMode() const override
std::optional< PublicKey > const validatorMasterPK_
void doTransactionSyncBatch(std::unique_lock< std::mutex > &lock, std::function< bool(std::unique_lock< std::mutex > const &)> retryCallback)
HashMap< std::uint64_t, InfoSub::Wptr > SubMapType
bool tryRemoveRpcSub(std::string const &strUrl) override
void doTransactionSync(std::shared_ptr< Transaction > transaction, bool bUnlimited, FailHard failType)
For transactions submitted directly by a client, apply batch of transactions and wait for this transa...
void submitTransaction(std::shared_ptr< STTx const > const &) override
bool subTransactions(InfoSub::Ref ispListener) override
void setAmendmentWarned() override
void pubPeerStatus(std::function< json::Value(void)> const &) override
void unsubMPT(InfoSub::Ref ispListener, HashSet< MPTID > const &mptIDs) override
StateAccounting accounting_
SubAccountHistoryMapType subAccountHistory_
SubMPTInfoMapType subMPT_
Guarded by mptLock_.
void setAccountHistoryJobTimer(SubAccountHistoryInfoWeak subInfo)
std::atomic< bool > unlBlocked_
void unsubAccountHistory(InfoSub::Ref ispListener, AccountID const &account, bool historyOnly) override
unsubscribe an account's transactions
bool subLedger(InfoSub::Ref ispListener, json::Value &jvResult) override
bool unsubBookChanges(std::uint64_t uListener) override
void unsubAccountInternal(std::uint64_t seq, HashSet< AccountID > const &vnaAccountIDs, bool rt) override
void setStateTimer() override
Called to initially start our timers.
std::size_t getLocalTxCount() override
bool preProcessTransaction(std::shared_ptr< Transaction > &transaction)
void processTransaction(std::shared_ptr< Transaction > &transaction, bool bUnlimited, bool bLocal, FailHard failType) override
Process transactions as they arrive from the network or which are submitted by clients.
RCLConsensus consensus_
bool unsubTransactions(std::uint64_t uListener) override
bool isAmendmentWarned() override
void getBookPage(std::shared_ptr< ReadView const > &lpLedger, Book const &, AccountID const &uTakerID, bool const bProof, unsigned int iLimit, json::Value const &jvMarker, json::Value &jvResult) override
std::mutex validationsMutex_
std::uint32_t acceptLedger(std::optional< std::chrono::milliseconds > consensusDelay) override
Accepts the current transaction tree, return the new ledger's sequence.
void clearUNLBlocked() override
void subAccount(InfoSub::Ref ispListener, HashSet< AccountID > const &vnaAccountIDs, bool rt) override
bool isUNLBlocked() override
void pubProposedTransaction(std::shared_ptr< ReadView const > const &ledger, std::shared_ptr< STTx const > const &transaction, TER result) override
std::atomic< bool > amendmentWarned_
boost::asio::steady_timer clusterTimer_
void cleanupSubscriptionMap(std::uint64_t seq, HashSet< AccountID > const &accounts, OuterMap &outerMap, BeforeErase &&beforeErase)
Erase one connection's entries from a subscription map in accountLock_-bounded chunks.
bool checkLastClosedLedger(Overlay::PeerSequence const &, UInt256 &networkClosed)
bool unsubPeerStatus(std::uint64_t uListener) override
void reportFeeChange() override
ServerFeeSummary lastFeeSummary_
void processTransactionSet(CanonicalTXSet const &set) override
Process a set of transactions synchronously, and ensuring that they are processed in one batch.
void mapComplete(std::shared_ptr< SHAMap > const &map, bool fromAcquire) override
bool isBlocked() override
~NetworkOPsImp() override
json::Value getServerInfo(bool human, bool admin, bool counters) override
Provides server functionality for clients.
Definition NetworkOPs.h:82
beast::AbstractClock< std::chrono::steady_clock > ClockType
Definition NetworkOPs.h:84
Writable ledger view that accumulates state and tx changes.
Definition OpenView.h:59
std::vector< std::shared_ptr< Peer > > PeerSequence
Definition Overlay.h:66
Manages the generic consensus algorithm for use by the RCL.
A peer's signed, proposed position for use in RCLConsensus.
PublicKey const & publicKey() const
Public key of peer that sent the proposal.
Represents a set of transactions in RCLConsensus.
Definition RCLCxTx.h:60
Wraps a ledger instance for use in generic Validations LedgerTrie.
static std::string getWordFromBlob(void const *blob, size_t bytes)
Chooses a single dictionary word from the data.
Definition RFC1751.cpp:434
Collects logging information.
std::unique_ptr< std::stringstream > const & ss()
A view into a ledger.
Definition ReadView.h:41
virtual SLE::const_pointer read(Keylet const &k) const =0
Return the state item associated with a key.
virtual std::optional< key_type > succ(key_type const &key, std::optional< key_type > const &last=std::nullopt) const =0
Return the key of the next state item.
std::vector< AccountTx > AccountTxs
constexpr bool holds() const noexcept
Definition STAmount.h:478
std::string getText() const override
Definition STAmount.cpp:647
Asset const & asset() const
Definition STAmount.h:496
void setJson(json::Value &) const
Definition STAmount.cpp:607
MPTAmount mpt() const
Definition STAmount.cpp:302
std::shared_ptr< STLedgerEntry > pointer
std::shared_ptr< STLedgerEntry const > const_pointer
Automatically unlocks and re-locks a unique_lock object.
Definition scope.h:197
std::size_t size() const noexcept
Definition Serializer.h:147
void const * data() const noexcept
Definition Serializer.h:153
Service registry for dependency injection.
std::flat_set< MPTID > getAffectedMPTs() const
Definition TxMeta.cpp:153
static time_point now()
Validator keys and manifest as set in configuration file.
json::Value jsonClipped() const
Definition XRPAmount.h:210
constexpr double decimalXRP() const
Definition XRPAmount.h:257
T duration_cast(T... args)
T emplace_back(T... args)
T emplace(T... args)
T empty(T... args)
T end(T... args)
T erase(T... args)
T exchange(T... args)
T find(T... args)
T get(T... args)
T insert(T... args)
T is_sorted(T... args)
T lock(T... args)
T make_pair(T... args)
T make_shared(T... args)
T make_unique(T... args)
T max(T... args)
T min(T... args)
void rngfill(void *const buffer, std::size_t const bytes, Generator &g)
Definition rngfill.h:11
constexpr Zero kZero
Definition Zero.h:30
JSON (JavaScript Object Notation).
Definition json_errors.h:5
int Int
unsigned int UInt
@ Array
array value (ordered list)
Definition json_value.h:28
@ Object
object value (collection of name/value pairs).
Definition json_value.h:29
STL namespace.
std::string const & getVersionString()
Server version.
Definition BuildInfo.cpp:63
TER valid(STTx const &tx, ReadView const &view, AccountID const &src, beast::Journal j)
std::string const & getCommitHash()
Definition Git.cpp:18
std::string const & getBuildBranch()
Definition Git.cpp:25
Keylet offer(AccountID const &id, SeqProxy const &seq) noexcept
An offer from an account.
Definition Indexes.cpp:298
Keylet book(Book const &b)
The beginning of an order book.
Definition Indexes.cpp:269
Keylet ownerDir(AccountID const &id) noexcept
The root page of an account's directory.
Definition Indexes.cpp:403
Keylet page(UInt256 const &root, std::uint64_t const index=0) noexcept
A page in a directory.
Definition Indexes.cpp:409
Keylet account(AccountID const &id) noexcept
AccountID root.
Definition Indexes.cpp:220
Keylet child(UInt256 const &key) noexcept
Any item that can be in an owner dir.
Definition Indexes.cpp:226
std::optional< std::string > encodeCTID(uint32_t ledgerSeq, uint32_t txnIndex, uint32_t networkID) noexcept
Encodes ledger sequence, transaction index, and network ID into a CTID string.
Definition CTID.h:36
static constexpr std::integral_constant< unsigned, Version > kApiVersion
Definition ApiVersion.h:39
void insertAllSyntheticInJson(json::Value &metadata, ReadView const &ledger, std::shared_ptr< STTx const > const &transaction, TxMeta const &transactionMeta)
Adds all synthetic fields to transaction metadata JSON.
json::Value computeBookChanges(std::shared_ptr< L const > const &lpAccepted)
Definition BookChanges.h:41
Rate rate(Env &env, Account const &account, std::uint32_t const &seq)
Definition escrow.cpp:60
Use hash_* containers for keys that do not need a cryptographically secure hashing algorithm.
Definition algorithm.h:5
bool cdirFirst(ReadView const &view, UInt256 const &root, SLE::const_pointer &page, unsigned int &index, UInt256 &entry)
Returns the first entry in the directory, advancing the index.
@ WarnRpcAmendmentBlocked
Definition ErrorCodes.h:159
@ WarnRpcUnsupportedMajority
Definition ErrorCodes.h:158
@ WarnRpcExpiredValidatorList
Definition ErrorCodes.h:160
STAmount divide(STAmount const &amount, Rate const &rate)
Definition Rate2.cpp:69
@ terQUEUED
Definition TER.h:226
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,...
constexpr FlagValue tfInnerBatchTxn
Definition TxFlags.h:44
bool isTerRetry(TER x) noexcept
Definition TER.h:677
ErrorCodeI
Definition ErrorCodes.h:23
@ RpcSuccess
Definition ErrorCodes.h:27
@ RpcInvalidParams
Definition ErrorCodes.h:67
void handleNewValidation(Application &app, std::shared_ptr< STValidation > const &val, std::string const &source, BypassAccept const bypassAccept, std::optional< beast::Journal > j)
Handle a new validation.
UInt256 getBookBase(Book const &book)
Definition Indexes.cpp:123
@ WrongLedger
We have the wrong ledger and are attempting to acquire it.
@ Proposing
We are normal participant in consensus and propose our position.
std::optional< std::uint64_t > mulDiv(std::uint64_t value, std::uint64_t mul, std::uint64_t div)
Return value*mul/div accurately.
SendIfPred< Predicate > sendIf(std::shared_ptr< Message > const &m, Predicate const &f)
Helper function to aid in type deduction.
Definition predicates.h:61
void forAllApiVersions(Fn const &fn, Args &&... args)
Definition ApiVersion.h:159
@ SigBad
Signature is bad.
Definition apply.h:28
@ Valid
Signature and local checks are good / passed.
Definition apply.h:36
T get(Section const &section, std::string const &name, T const &defaultValue=T{})
Retrieve a key/value pair from a section.
constexpr std::size_t kMaxPoppedTransactions
std::string strHex(FwdIt begin, FwdIt end)
Definition strHex.h:13
std::uint64_t getQuality(UInt256 const &uBase)
Definition Indexes.cpp:194
std::unique_ptr< FeeVote > makeFeeVote(FeeSetup const &setup, beast::Journal journal)
Create an instance of the FeeVote logic.
std::pair< Validity, std::string > checkValidity(HashRouter &router, STTx const &tx, Rules const &rules)
Checks transaction signature and local checks.
Definition apply.cpp:62
Rules makeRulesGivenLedger(DigestAwareReadView const &ledger, Rules const &current)
Definition ReadView.cpp:62
FeeSetup setupFeeVote(Section const &section)
@ tefPAST_SEQ
Definition TER.h:170
std::string toBase58(AccountID const &v)
Convert AccountID to base58 checked string.
Definition AccountID.cpp:95
std::unique_ptr< NetworkOPs > makeNetworkOPs(ServiceRegistry &registry, NetworkOPs::ClockType &clock, bool standalone, std::size_t minPeerCount, bool startValid, JobQueue &jobQueue, LedgerMaster &ledgerMaster, ValidatorKeys const &validatorKeys, boost::asio::io_context &ioSvc, beast::Journal journal, beast::insight::Collector::Ptr const &collector)
UInt256 getQualityNext(UInt256 const &uBase)
Definition Indexes.cpp:186
Number root(Number f, unsigned d)
CsprngEngine & cryptoPrng()
The default cryptographically secure PRNG.
bool transResultInfo(TER code, std::string &token, std::string &text)
Definition TER.cpp:242
Seed generateSeed(std::string const &passPhrase)
Generate a seed deterministically.
Definition Seed.cpp:58
constexpr Dest safeCast(Src s) noexcept
Definition safe_cast.h:21
std::string transToken(TER code)
Definition TER.cpp:257
std::pair< PublicKey, SecretKey > generateKeyPair(KeyType type, Seed const &seed)
Generate a key pair deterministically.
constexpr std::uint32_t kFeeUnitsDeprecated
std::string to_string(BaseUInt< Bits, Tag > const &a)
Definition base_uint.h:657
std::string toStringIso(date::sys_time< Duration > tp)
Definition chrono.h:70
STAmount accountFunds(ReadView const &view, AccountID const &id, STAmount const &saDefault, FreezeHandling freezeHandling, beast::Journal j)
bool isGlobalFrozen(ReadView const &view, AccountID const &issuer)
Check if the issuer has the global freeze flag set.
BaseUInt< 256 > UInt256
Definition base_uint.h:580
STAmount amountFromQuality(std::uint64_t rate)
Definition STAmount.cpp:896
@ JtClientAcctHist
Definition Job.h:35
@ JtTransaction
Definition Job.h:48
@ JtBatch
Definition Job.h:51
@ JtTxnProc
Definition Job.h:68
@ JtClientConsensus
Definition Job.h:34
@ JtNetopCluster
Definition Job.h:61
@ JtClientFeeChange
Definition Job.h:33
bool isTefFailure(TER x) noexcept
Definition TER.h:671
HashSet< Book > affectedBooks(AcceptedLedgerTx const &alTx, beast::Journal const &j)
Extract the set of books affected by a transaction.
std::unique_ptr< LocalTxs > makeLocalTxs()
Definition LocalTxs.cpp:185
IOUAmount mulRatio(IOUAmount const &amt, std::uint32_t num, std::uint32_t den, bool roundUp)
Rate transferRate(ReadView const &view, AccountID const &issuer)
Returns IOU issuer transfer fee as Rate.
BaseUInt< 192 > MPTID
MPTID is a 192-bit value representing MPT Issuance ID, which is a concatenation of a 32-bit sequence ...
Definition UintTypes.h:54
ConsensusPhase
Phases of consensus for a single ledger round.
@ Open
We haven't closed our ledger yet, but others might have.
Rate const kParityRate
A transfer rate signifying a 1:1 exchange.
json::Value getJson(LedgerFill const &fill)
Return a new json::Value representing the ledger with given options.
std::unordered_set< Value, Hash, Pred, Allocator > HashSet
bool cdirNext(ReadView const &view, UInt256 const &root, SLE::const_pointer &page, unsigned int &index, UInt256 &entry)
Returns the next entry in the directory, advancing the index.
AccountID calcAccountID(PublicKey const &pk)
static std::array< char const *, 5 > const kStateNames
constexpr auto kMuldivMax
Definition mulDiv.h:8
ApplyFlags
Definition ApplyView.h:27
@ TapUnlimited
Definition ApplyView.h:39
@ TapFailHard
Definition ApplyView.h:32
@ TapNone
Definition ApplyView.h:28
BaseUInt< 160, detail::AccountIDTag > AccountID
A 160-bit unsigned that uniquely identifies an account.
Definition AccountID.h:34
detail::MultiApiJson< rpc::kApiMinimumSupportedVersion, rpc::kApiMaximumValidVersion > MultiApiJson
bool isTelLocal(TER x) noexcept
Definition TER.h:659
@ temINVALID_FLAG
Definition TER.h:99
@ temBAD_SIGNATURE
Definition TER.h:93
bool isTesSuccess(TER x) noexcept
Definition TER.h:683
static std::uint32_t trunc32(std::uint64_t v)
TERSubset< CanCvtToTER > TER
Definition TER.h:654
STAmount multiply(STAmount const &amount, Number const &frac, Number::RoundingMode rm)
std::unordered_map< Key, Value, Hash, Pred, Allocator > HashMap
OperatingMode
Specifies the mode under which the server believes it's operating.
Definition NetworkOPs.h:60
@ TRACKING
convinced we agree with the network
Definition NetworkOPs.h:64
@ DISCONNECTED
not ready to process requests
Definition NetworkOPs.h:61
@ CONNECTED
convinced we are talking to the network
Definition NetworkOPs.h:62
@ FULL
we have the ledger and can even validate
Definition NetworkOPs.h:65
@ SYNCING
fallen slightly behind
Definition NetworkOPs.h:63
std::shared_ptr< STTx const > sterilize(STTx const &stx)
Sterilize a transaction.
Definition STTx.cpp:880
STAmount issuerFundsToSelfIssue(ReadView const &view, MPTIssue const &issue)
Determine funds available for an issuer to sell in an issuer owned offer.
static auto const kGenesisAccountId
bool isTemMalformed(TER x) noexcept
Definition TER.h:665
STAmount accountHolds(ReadView const &view, AccountID const &account, Currency const &currency, AccountID const &issuer, FreezeHandling zeroIfFrozen, beast::Journal j, SpendableHandling includeFullBalance=SpendableHandling::SimpleBalance)
@ tesSUCCESS
Definition TER.h:250
XRPL_NO_SANITIZE_ADDRESS void Throw(Args &&... args)
Definition contract.h:52
STAmount toSTAmount(IOUAmount const &iou, Asset const &asset)
CanonicalTXSet OrderedTxs
Definition OpenLedger.h:35
T owns_lock(T... args)
T push_back(T... args)
T ref(T... args)
T reserve(T... args)
T reset(T... args)
T set_intersection(T... args)
T size(T... args)
T str(T... args)
static constexpr auto kPort
Definition Constants.h:146
static constexpr auto kIp
Definition Constants.h:114
PublicKey masterKey
The master key associated with this manifest.
Definition Manifest.h:84
std::string serialized
The manifest in serialized form.
Definition Manifest.h:79
Blob getMasterSignature() const
Returns manifest master key signature.
std::string domain
The domain, if one was specified in the manifest; empty otherwise.
Definition Manifest.h:102
std::optional< Blob > getSignature() const
Returns manifest signature.
std::optional< PublicKey > signingKey
The ephemeral key associated with this manifest.
Definition Manifest.h:92
std::uint32_t sequence
The sequence number of this manifest.
Definition Manifest.h:97
Server fees published on server subscription.
std::optional< TxQ::Metrics > em
bool operator!=(ServerFeeSummary const &b) const
bool operator==(ServerFeeSummary const &b) const
beast::insight::Gauge fullTransitions
beast::insight::Gauge disconnectedTransitions
beast::insight::Gauge connectedDuration
beast::insight::Gauge trackingTransitions
beast::insight::Gauge fullDuration
beast::insight::Gauge syncingDuration
Stats(Handler const &handler, beast::insight::Collector::Ptr const &collector)
beast::insight::Gauge disconnectedDuration
beast::insight::Gauge connectedTransitions
beast::insight::Gauge trackingDuration
beast::insight::Hook hook
beast::insight::Gauge syncingTransitions
SubAccountHistoryIndex(AccountID const &accountId)
std::shared_ptr< SubAccountHistoryIndex > index
std::shared_ptr< SubAccountHistoryIndex > index
Select all peers (except optional excluded) that are in our cluster.
Definition predicates.h:127
Represents a transfer rate.
Definition Rate.h:21
std::uint32_t value
Definition Rate.h:22
static constexpr auto kPortGrpc
Definition Constants.h:45
Sends a message to all peers.
Definition predicates.h:15
Changes in trusted nodes after updating validator list.
HashSet< NodeID > added
HashSet< NodeID > removed
Structure returned by TxQ::getMetrics, expressed in reference fee level units.
Definition TxQ.h:186
void set(char const *key, auto const &v)
IsMemberResult isMember(char const *key) const
Data format for exchanging consumption information across peers.
Definition Gossip.h:13
std::vector< Item > items
Definition Gossip.h:27
T time_since_epoch(T... args)
T to_string(T... args)
T unlock(T... args)
T value_or(T... args)
T what(T... args)