xrpld
Loading...
Searching...
No Matches
libxrpl/server/InfoSub.cpp
1#include <xrpl/server/InfoSub.h>
2
3#include <xrpl/basics/Log.h>
4#include <xrpl/basics/UnorderedContainers.h>
5#include <xrpl/beast/utility/Journal.h>
6#include <xrpl/beast/utility/instrumentation.h>
7#include <xrpl/protocol/AccountID.h>
8#include <xrpl/protocol/Book.h>
9#include <xrpl/protocol/UintTypes.h>
10#include <xrpl/resource/Consumer.h>
11
12#include <cstddef>
13#include <cstdint>
14#include <exception>
15#include <memory>
16#include <mutex>
17#include <utility>
18
19namespace xrpl {
20
21namespace {
22
23// Wraps a Source teardown call so that an exception from one cleanup
24// step does not prevent the subsequent steps from running. Source methods
25// acquire a lock and can throw std::system_error; a throw out of ~InfoSub
26// during stack unwinding would terminate the process. Failures are
27// reported through the Source's Journal so they reach the configured log
28// sinks; JLOG itself cannot throw, so the noexcept guarantee holds.
29template <typename F>
30void
31safeUnsub(std::uint64_t seq, F&& f, beast::Journal j) noexcept
32{
33 try
34 {
35 f();
36 }
37 catch (std::exception const& e)
38 {
39 JLOG(j.warn()) << "~InfoSub[seq=" << seq << "]: cleanup step failed: " << e.what();
40 }
41 catch (...)
42 {
43 JLOG(j.warn()) << "~InfoSub[seq=" << seq << "]: cleanup step failed: unknown exception";
44 }
45}
46
47// Number of requested entries not already tracked in `existing`. Only these are
48// charged against the cap, so re-subscribing entries a connection already holds
49// is free.
50template <typename T>
51[[nodiscard]] std::size_t
52countNew(HashSet<T> const& requested, HashSet<T> const& existing)
53{
54 std::size_t fresh = 0;
55 for (auto const& entry : requested)
56 {
57 if (!existing.contains(entry))
58 ++fresh;
59 }
60 return fresh;
61}
62
63} // namespace
64
65// This is the primary interface into the "client" portion of the program.
66// Code that wants to do normal operations on the network such as
67// creating and monitoring accounts, creating transactions, and so on
68// should use this interface. The RPC code will primarily be a light wrapper
69// over this code.
70
71// Eventually, it will check the node's operating mode (synced, unsynced,
72// etcetera) and defer to the correct means of processing. The current
73// code assumes this node is synced (and will continue to do so until
74// there's a functional network.
75
77{
78}
79
81 : consumer_(consumer), source_(source), seq_(assignId())
82{
83}
84
86{
87 // Stream unsubscribes are O(1): each erases this connection's single seq_
88 // from one stream map, so they are cheap enough to run inline on the
89 // disconnect thread.
90 // Each Source teardown call below acquires a server-side lock and
91 // can throw. Wrap each independent call so partial failure does not
92 // skip the remaining teardown steps.
93
94 auto const& j = source_.journal();
95
96 safeUnsub(seq_, [&] { source_.unsubTransactions(seq_); }, j);
97 safeUnsub(seq_, [&] { source_.unsubRTTransactions(seq_); }, j);
98 safeUnsub(seq_, [&] { source_.unsubLedger(seq_); }, j);
99 safeUnsub(seq_, [&] { source_.unsubManifests(seq_); }, j);
100 safeUnsub(seq_, [&] { source_.unsubServer(seq_); }, j);
101 safeUnsub(seq_, [&] { source_.unsubValidations(seq_); }, j);
102 safeUnsub(seq_, [&] { source_.unsubPeerStatus(seq_); }, j);
103 safeUnsub(seq_, [&] { source_.unsubConsensus(seq_); }, j);
104
105 // MPT subscriptions are torn down inline here, keyed on seq_, like books
106 // below. The set is capped; each unsubMPTInternal takes mptLock_ for a
107 // single O(1) erase and releases it, so a competing MPT publish can
108 // interleave between erases. Use the internal variant so it does not write
109 // back to mptSubscriptions_ on this partially-destroyed object.
110 for (auto const& mptID : mptSubscriptions_)
111 {
112 safeUnsub(seq_, [&] { source_.unsubMPTInternal(seq_, mptID); }, j);
113 }
114
115 // Book subscriptions are torn down inline here, keyed on seq_, rather than
116 // through the chunked account cleanup below. The book set is not capped, so
117 // it can be large; but each unsubBookInternal takes bookLock_ for a single
118 // O(1) erase and releases it, so even a large set never holds a lock across
119 // the whole loop - a competing book publish can interleave between erases.
120 // The disconnect thread still does O(N) brief acquisitions. Use the internal
121 // variant so it does not write back to bookSubscriptions_ on this
122 // partially-destroyed object.
123 for (auto const& book : bookSubscriptions_)
124 {
125 safeUnsub(seq_, [&] { source_.unsubBookInternal(seq_, book); }, j);
126 }
127
128 // Hand the account sets off (by move) to the Source for a chunked,
129 // off-thread teardown keyed on seq_, instead of erasing them inline here.
130 // This keeps the destructor from holding the account lock across a large
131 // erase loop. The job never references this object, which is being
132 // destroyed.
133 //
134 // Moving the sets without holding lock_ is safe: the destructor runs only
135 // when the last shared_ptr to this InfoSub is released, so by the
136 // shared_ptr contract no other thread holds a reference. Subscription maps
137 // store weak_ptrs, so a concurrent publisher must weak_ptr::lock() first;
138 // that succeeds only while a strong reference exists, which cannot overlap
139 // with destruction. No other thread can observe the moved-from sets.
140 //
141 // Wrapped like the steps above: scheduleAccountCleanup enqueues a JobQueue
142 // task, which allocates and locks and so can throw. A throw out of this
143 // noexcept destructor would terminate the process. Skipping the cleanup on
144 // throw is harmless: the account/rt maps hold weak_ptrs that the next
145 // publish prunes once this InfoSub is gone, and any history paging job
146 // self-terminates when its weak sink can no longer be locked.
147 safeUnsub(
148 seq_,
149 [&] {
150 source_.scheduleAccountCleanup(
151 seq_,
152 std::move(realTimeSubscriptions_),
153 std::move(normalSubscriptions_),
155 },
156 j);
157}
158
161{
162 return consumer_;
163}
164
167{
168 return seq_;
169}
170
171void
175
182
185{
186 // Hold lock_ for the whole read so the counted sets cannot be mutated
187 // mid-count by a concurrent (un)subscribe on this connection.
188 std::scoped_lock const sl(lock_);
189
190 return subscriptionCount(sl);
191}
192
193bool
195 HashSet<AccountID> const& proposedAccounts,
196 HashSet<AccountID> const& normalAccounts,
197 std::size_t cap)
198{
199 // One lock hold covers the count, the check and the insert.
200 std::scoped_lock const sl(lock_);
201
202 std::size_t const additional = countNew(proposedAccounts, realTimeSubscriptions_) +
203 countNew(normalAccounts, normalSubscriptions_);
204
205 if (exceedsSubscriptionCap(subscriptionCount(sl), additional, cap))
206 return false;
207
208 realTimeSubscriptions_.insert(proposedAccounts.begin(), proposedAccounts.end());
209 normalSubscriptions_.insert(normalAccounts.begin(), normalAccounts.end());
210 return true;
211}
212
213bool
215{
216 // One lock hold covers the count, the check and the insert.
217 std::scoped_lock const sl(lock_);
218
219 if (exceedsSubscriptionCap(subscriptionCount(sl), countNew(mptIDs, mptSubscriptions_), cap))
220 return false;
221
222 mptSubscriptions_.insert(mptIDs.begin(), mptIDs.end());
223 return true;
224}
225
226void
228{
229 std::scoped_lock const sl(lock_);
230
231 if (rt)
232 {
233 realTimeSubscriptions_.insert(account);
234 }
235 else
236 {
237 normalSubscriptions_.insert(account);
238 }
239}
240
241void
243{
244 std::scoped_lock const sl(lock_);
245
246 if (rt)
247 {
248 realTimeSubscriptions_.erase(account);
249 }
250 else
251 {
252 normalSubscriptions_.erase(account);
253 }
254}
255
256bool
258{
259 std::scoped_lock const sl(lock_);
260 return accountHistorySubscriptions_.insert(account).second;
261}
262
263void
265{
266 std::scoped_lock const sl(lock_);
267 accountHistorySubscriptions_.erase(account);
268}
269
270bool
272{
273 std::scoped_lock const sl(lock_);
274 return accountHistorySubscriptions_.contains(account);
275}
276
277void
279{
280 std::scoped_lock const sl(lock_);
281 bookSubscriptions_.insert(book);
282}
283
284void
286{
287 std::scoped_lock const sl(lock_);
288 bookSubscriptions_.erase(book);
289}
290
291void
293{
294 request_.reset();
295}
296
297void
302
305{
306 return request_;
307}
308
309void
310InfoSub::setApiVersion(unsigned int apiVersion)
311{
312 apiVersion_ = apiVersion;
313}
314
315unsigned int
317{
318 XRPL_ASSERT(apiVersion_ > 0, "xrpl::InfoSub::getApiVersion : valid API version");
319 return apiVersion_;
320}
321
322void
324{
325 std::scoped_lock const sl(lock_);
326
327 mptSubscriptions_.insert(mptID);
328}
329
330void
332{
333 std::scoped_lock const sl(lock_);
334
335 mptSubscriptions_.erase(mptID);
336}
337
338} // namespace xrpl
T begin(T... args)
Specifies an order book.
Definition Book.h:28
Abstracts the source of subscription data.
Definition InfoSub.h:107
void setRequest(std::shared_ptr< InfoSubRequest > const &req)
void insertBookSubscription(Book const &book)
Record that this subscriber is following book.
InfoSub(Source &source)
Consumer consumer_
Definition InfoSub.h:479
bool insertSubAccountHistory(AccountID const &account)
bool tryReserveMPTSubscriptions(HashSet< MPTID > const &mptIDs, std::size_t cap)
Enforce the cap and reserve a request's net-new MPT issuances, atomically.
HashSet< MPTID > mptSubscriptions_
Definition InfoSub.h:487
bool tryReserveAccountSubscriptions(HashSet< AccountID > const &proposedAccounts, HashSet< AccountID > const &normalAccounts, std::size_t cap)
Enforce the cap and reserve a request's net-new accounts, atomically.
void deleteSubMPTInfo(MPTID const &mptID)
std::size_t totalSubscriptionCount() const
Return the number of subscriptions currently tracked on this connection.
void setApiVersion(unsigned int apiVersion)
static int assignId()
Definition InfoSub.h:491
std::uint64_t getSeq() const
HashSet< AccountID > realTimeSubscriptions_
Definition InfoSub.h:481
std::scoped_lock< decltype(lock_)> ScopedLock
Definition InfoSub.h:469
resource::Consumer Consumer
Definition InfoSub.h:100
void insertSubAccountInfo(AccountID const &account, bool rt)
std::size_t subscriptionCount(ScopedLock const &lock) const
The combined tally the per-connection cap is enforced against.
void deleteBookSubscription(Book const &book)
Stop tracking book for this subscriber.
bool hasAccountHistorySubscription(AccountID const &account) const
Whether this connection already tracks an account-history for account.
Source & source_
Definition InfoSub.h:480
HashSet< AccountID > normalSubscriptions_
Definition InfoSub.h:482
HashSet< Book > bookSubscriptions_
Definition InfoSub.h:486
void insertSubMPTInfo(MPTID const &mptID)
HashSet< AccountID > accountHistorySubscriptions_
Definition InfoSub.h:485
void deleteSubAccountInfo(AccountID const &account, bool rt)
void deleteSubAccountHistory(AccountID const &account)
unsigned int getApiVersion() const noexcept
std::uint64_t seq_
Definition InfoSub.h:484
std::mutex lock_
Definition InfoSub.h:465
unsigned int apiVersion_
Definition InfoSub.h:488
std::shared_ptr< InfoSubRequest > const & getRequest()
std::shared_ptr< InfoSubRequest > request_
Definition InfoSub.h:483
An endpoint that consumes resources.
Definition Consumer.h:20
T end(T... args)
Use hash_* containers for keys that do not need a cryptographically secure hashing algorithm.
Definition algorithm.h:5
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
std::unordered_set< Value, Hash, Pred, Allocator > HashSet
BaseUInt< 160, detail::AccountIDTag > AccountID
A 160-bit unsigned that uniquely identifies an account.
Definition AccountID.h:34
constexpr bool exceedsSubscriptionCap(std::size_t current, std::size_t additional, std::size_t cap=kMaxSubscriptionsPerConnection)
Whether adding additional subscriptions to a connection already holding current would exceed the cap.
Definition InfoSub.h:52
T what(T... args)