xrpld
Loading...
Searching...
No Matches
resource/detail/Logic.h
1#pragma once
2
3#include <xrpl/basics/Log.h>
4#include <xrpl/basics/UnorderedContainers.h>
5#include <xrpl/basics/chrono.h>
6#include <xrpl/beast/clock/abstract_clock.h>
7#include <xrpl/beast/core/List.h>
8#include <xrpl/beast/insight/Collector.h>
9#include <xrpl/beast/net/IPEndpoint.h>
10#include <xrpl/beast/utility/Journal.h>
11#include <xrpl/beast/utility/PropertyStream.h>
12#include <xrpl/beast/utility/instrumentation.h>
13#include <xrpl/json/json_value.h>
14#include <xrpl/protocol/jss.h>
15#include <xrpl/resource/Charge.h>
16#include <xrpl/resource/Consumer.h>
17#include <xrpl/resource/Disposition.h>
18#include <xrpl/resource/Fees.h>
19#include <xrpl/resource/Gossip.h>
20#include <xrpl/resource/detail/Import.h>
21
22#include <mutex>
23#include <string>
24#include <tuple>
25#include <utility>
26
27namespace xrpl::resource {
28
29class Logic
30{
31private:
36
37 struct Stats
38 {
40 {
41 warn = collector->makeMeter("warn");
42 drop = collector->makeMeter("drop");
43 }
44
47 };
48
52
54
55 // Table of all entries
57
58 // Because the following are intrusive lists, a given Entry may be in
59 // at most list at a given instant. The Entry must be removed from
60 // one list before placing it in another.
61
62 // List of all active inbound entries
64
65 // List of all active outbound entries
67
68 // List of all active admin entries
70
71 // List of all inactive entries
73
74 // All imported gossip data
76
77 //--------------------------------------------------------------------------
78public:
80 : stats_(collector), clock_(clock), journal_(journal)
81 {
82 }
83
85 {
86 // These have to be cleared before the Logic is destroyed
87 // since their destructors call back into the class.
88 // Order matters here as well, the import table has to be
89 // destroyed before the consumer table.
90 //
91 importTable_.clear();
92 table_.clear();
93 }
94
97 {
98 Entry* entry(nullptr);
99
100 {
101 std::scoped_lock const _(lock_);
102 auto [resultIt, resultInserted] = table_.emplace(
104 std::make_tuple(Kind::Inbound, address.atPort(0)), // Key
105 std::make_tuple(clock_.now())); // Entry
106
107 entry = &resultIt->second;
108 entry->key = &resultIt->first;
109 ++entry->refcount;
110 if (entry->refcount == 1)
111 {
112 if (!resultInserted)
113 {
114 inactive_.erase(inactive_.iteratorTo(*entry));
115 }
116 inbound_.pushBack(*entry);
117 }
118 }
119
120 JLOG(journal_.debug()) << "New inbound endpoint " << *entry;
121
122 return Consumer(*this, *entry);
123 }
124
127 {
128 Entry* entry(nullptr);
129
130 {
131 std::scoped_lock const _(lock_);
132 auto [resultIt, resultInserted] = table_.emplace(
134 std::make_tuple(Kind::Outbound, address), // Key
135 std::make_tuple(clock_.now())); // Entry
136
137 entry = &resultIt->second;
138 entry->key = &resultIt->first;
139 ++entry->refcount;
140 if (entry->refcount == 1)
141 {
142 if (!resultInserted)
143 inactive_.erase(inactive_.iteratorTo(*entry));
144 outbound_.pushBack(*entry);
145 }
146 }
147
148 JLOG(journal_.debug()) << "New outbound endpoint " << *entry;
149
150 return Consumer(*this, *entry);
151 }
152
160 {
161 Entry* entry(nullptr);
162
163 {
164 std::scoped_lock const _(lock_);
165 auto [resultIt, resultInserted] = table_.emplace(
167 std::make_tuple(Kind::Unlimited, address.atPort(1)), // Key
168 std::make_tuple(clock_.now())); // Entry
169
170 entry = &resultIt->second;
171 entry->key = &resultIt->first;
172 ++entry->refcount;
173 if (entry->refcount == 1)
174 {
175 if (!resultInserted)
176 inactive_.erase(inactive_.iteratorTo(*entry));
177 admin_.pushBack(*entry);
178 }
179 }
180
181 JLOG(journal_.debug()) << "New unlimited endpoint " << *entry;
182
183 return Consumer(*this, *entry);
184 }
185
188 {
190 }
191
196 getJson(int threshold)
197 {
198 ClockType::time_point const now(clock_.now());
199
201 std::scoped_lock const _(lock_);
202
203 for (auto& inboundEntry : inbound_)
204 {
205 int const localBalance = inboundEntry.localBalance.value(now);
206 if ((localBalance + inboundEntry.remoteBalance) >= threshold)
207 {
208 json::Value& entry = (ret[inboundEntry.toString()] = json::ValueType::Object);
209 entry[jss::local] = localBalance;
210 entry[jss::remote] = inboundEntry.remoteBalance;
211 entry[jss::type] = "inbound";
212 }
213 }
214 for (auto& outboundEntry : outbound_)
215 {
216 int const localBalance = outboundEntry.localBalance.value(now);
217 if ((localBalance + outboundEntry.remoteBalance) >= threshold)
218 {
219 json::Value& entry = (ret[outboundEntry.toString()] = json::ValueType::Object);
220 entry[jss::local] = localBalance;
221 entry[jss::remote] = outboundEntry.remoteBalance;
222 entry[jss::type] = "outbound";
223 }
224 }
225 for (auto& adminEntry : admin_)
226 {
227 int const localBalance = adminEntry.localBalance.value(now);
228 if ((localBalance + adminEntry.remoteBalance) >= threshold)
229 {
230 json::Value& entry = (ret[adminEntry.toString()] = json::ValueType::Object);
231 entry[jss::local] = localBalance;
232 entry[jss::remote] = adminEntry.remoteBalance;
233 entry[jss::type] = "admin";
234 }
235 }
236
237 return ret;
238 }
239
240 Gossip
242 {
243 ClockType::time_point const now(clock_.now());
244
245 Gossip gossip;
246 std::scoped_lock const _(lock_);
247
248 gossip.items.reserve(inbound_.size());
249
250 for (auto& inboundEntry : inbound_)
251 {
252 Gossip::Item item;
253 item.balance = inboundEntry.localBalance.value(now);
254 if (item.balance >= kMinimumGossipBalance)
255 {
256 item.address = inboundEntry.key->address;
257 gossip.items.push_back(item);
258 }
259 }
260
261 return gossip;
262 }
263
264 //--------------------------------------------------------------------------
265
266 void
267 importConsumers(std::string const& origin, Gossip const& gossip)
268 {
269 auto const elapsed = clock_.now();
270 {
271 std::scoped_lock const _(lock_);
272 auto [resultIt, resultInserted] = importTable_.emplace(
274 std::make_tuple(origin), // Key
275 std::make_tuple(clock_.now().time_since_epoch().count())); // Import
276
277 if (resultInserted)
278 {
279 // This is a new import
280 Import& next(resultIt->second);
281 next.whenExpires = elapsed + kGossipExpirationSeconds;
282 next.items.reserve(gossip.items.size());
283
284 for (auto const& gossipItem : gossip.items)
285 {
286 Import::Item item;
287 item.balance = gossipItem.balance;
288 item.consumer = newInboundEndpoint(gossipItem.address);
289 item.consumer.entry().remoteBalance += item.balance;
290 next.items.push_back(item);
291 }
292 }
293 else
294 {
295 // Previous import exists so add the new remote
296 // balances and then deduct the old remote balances.
297
298 Import next;
299 next.whenExpires = elapsed + kGossipExpirationSeconds;
300 next.items.reserve(gossip.items.size());
301 for (auto const& gossipItem : gossip.items)
302 {
303 Import::Item item;
304 item.balance = gossipItem.balance;
305 item.consumer = newInboundEndpoint(gossipItem.address);
306 item.consumer.entry().remoteBalance += item.balance;
307 next.items.push_back(item);
308 }
309
310 Import& prev(resultIt->second);
311 for (auto& item : prev.items)
312 {
313 item.consumer.entry().remoteBalance -= item.balance;
314 }
315
316 std::swap(next, prev);
317 }
318 }
319 }
320
321 //--------------------------------------------------------------------------
322
323 // Called periodically to expire entries and groom the table.
324 //
325 void
327 {
328 std::scoped_lock const _(lock_);
329
330 auto const elapsed = clock_.now();
331
332 for (auto iter(inactive_.begin()); iter != inactive_.end();)
333 {
334 if (iter->whenExpires <= elapsed)
335 {
336 JLOG(journal_.debug()) << "Expired " << *iter;
337 auto tableIter = table_.find(*iter->key);
338 ++iter;
339 erase(tableIter);
340 }
341 else
342 {
343 break;
344 }
345 }
346
347 auto iter = importTable_.begin();
348 while (iter != importTable_.end())
349 {
350 Import& import(iter->second);
351 if (iter->second.whenExpires <= elapsed)
352 {
353 for (auto& item : import.items)
354 {
355 item.consumer.entry().remoteBalance -= item.balance;
356 }
357
358 iter = importTable_.erase(iter);
359 }
360 else
361 {
362 ++iter;
363 }
364 }
365 }
366
367 //--------------------------------------------------------------------------
368
369 // Returns the disposition based on the balance and thresholds
370 static Disposition
372 {
373 if (balance >= kDropThreshold)
374 return Disposition::Drop;
375
377 return Disposition::Warn;
378
379 return Disposition::Ok;
380 }
381
382 void
383 erase(Table::iterator iter)
384 {
385 std::scoped_lock const _(lock_);
386 Entry& entry(iter->second);
387 XRPL_ASSERT(entry.refcount == 0, "xrpl::resource::Logic::erase : entry not used");
388 inactive_.erase(inactive_.iteratorTo(entry));
389 table_.erase(iter);
390 }
391
392 void
394 {
395 std::scoped_lock const _(lock_);
396 ++entry.refcount;
397 }
398
399 void
401 {
402 std::scoped_lock const _(lock_);
403 if (--entry.refcount == 0)
404 {
405 JLOG(journal_.debug()) << "Inactive " << entry;
406
407 switch (entry.key->kind)
408 {
409 case Kind::Inbound:
410 inbound_.erase(inbound_.iteratorTo(entry));
411 break;
412 case Kind::Outbound:
413 outbound_.erase(outbound_.iteratorTo(entry));
414 break;
415 case Kind::Unlimited:
416 admin_.erase(admin_.iteratorTo(entry));
417 break;
418 default:
419 // LCOV_EXCL_START
420 UNREACHABLE(
421 "xrpl::resource::Logic::release : invalid entry "
422 "kind");
423 break;
424 // LCOV_EXCL_STOP
425 }
426 inactive_.pushBack(entry);
427 entry.whenExpires = clock_.now() + kSecondsUntilExpiration;
428 }
429 }
430
432 charge(Entry& entry, Charge const& fee, std::string context = {})
433 {
434 static constexpr Charge::value_type kFeeLogAsWarn = 3000;
435 static constexpr Charge::value_type kFeeLogAsInfo = 1000;
436 static constexpr Charge::value_type kFeeLogAsDebug = 100;
437 static_assert(
438 kFeeLogAsWarn > kFeeLogAsInfo && kFeeLogAsInfo > kFeeLogAsDebug && kFeeLogAsDebug > 10);
439
440 static auto kGetStream = [](resource::Charge::value_type cost, beast::Journal& journal) {
441 if (cost >= kFeeLogAsWarn)
442 return journal.warn();
443 if (cost >= kFeeLogAsInfo)
444 return journal.info();
445 if (cost >= kFeeLogAsDebug)
446 return journal.debug();
447 return journal.trace();
448 };
449
450 if (!context.empty())
451 context = " (" + context + ")";
452
453 std::scoped_lock const _(lock_);
454 ClockType::time_point const now(clock_.now());
455 int const balance(entry.add(fee.cost(), now));
456 JLOG(kGetStream(fee.cost(), journal_)) << "Charging " << entry << " for " << fee << context;
457 return disposition(balance);
458 }
459
460 bool
461 warn(Entry& entry)
462 {
463 if (entry.isUnlimited())
464 return false;
465
466 std::scoped_lock const _(lock_);
467 bool notify(false);
468 auto const elapsed = clock_.now();
469 if (entry.balance(clock_.now()) >= kWarningThreshold && elapsed != entry.lastWarningTime)
470 {
471 charge(entry, kFeeWarning);
472 notify = true;
473 entry.lastWarningTime = elapsed;
474 }
475 if (notify)
476 {
477 JLOG(journal_.info()) << "Load warning: " << entry;
478 ++stats_.warn;
479 }
480 return notify;
481 }
482
483 bool
485 {
486 if (entry.isUnlimited())
487 return false;
488
489 std::scoped_lock const _(lock_);
490 bool drop(false);
491 ClockType::time_point const now(clock_.now());
492 int const balance(entry.balance(now));
493 if (balance >= kDropThreshold)
494 {
495 JLOG(journal_.warn()) << "Consumer entry " << entry << " dropped with balance "
496 << balance << " at or above drop threshold " << kDropThreshold;
497
498 // Adding feeDrop at this point keeps the dropped connection
499 // from re-connecting for at least a little while after it is
500 // dropped.
501 charge(entry, kFeeDrop);
502 ++stats_.drop;
503 drop = true;
504 }
505 return drop;
506 }
507
508 int
510 {
511 std::scoped_lock const _(lock_);
512 return entry.balance(clock_.now());
513 }
514
515 //--------------------------------------------------------------------------
516
517 static void
519 ClockType::time_point const now,
521 EntryIntrusiveList& list)
522 {
523 for (auto& entry : list)
524 {
525 beast::PropertyStream::Map item(items);
526 if (entry.refcount != 0)
527 item["count"] = entry.refcount;
528 item["name"] = entry.toString();
529 item["balance"] = entry.balance(now);
530 if (entry.remoteBalance != 0)
531 item["remote_balance"] = entry.remoteBalance;
532 }
533 }
534
535 void
537 {
538 ClockType::time_point const now(clock_.now());
539
540 std::scoped_lock const _(lock_);
541
542 {
543 beast::PropertyStream::Set s("inbound", map);
544 writeList(now, s, inbound_);
545 }
546
547 {
548 beast::PropertyStream::Set s("outbound", map);
549 writeList(now, s, outbound_);
550 }
551
552 {
553 beast::PropertyStream::Set s("admin", map);
554 writeList(now, s, admin_);
555 }
556
557 {
558 beast::PropertyStream::Set s("inactive", map);
559 writeList(now, s, inactive_);
560 }
561 }
562};
563
564} // namespace xrpl::resource
std::chrono::steady_clock::time_point time_point
virtual time_point now() const =0
Returns the current time.
A generic endpoint for log messages.
Definition Journal.h:44
Intrusive doubly linked list.
Definition List.h:258
std::shared_ptr< Collector > Ptr
Definition Collector.h:29
A metric for measuring an integral value.
Definition Meter.h:19
A version-independent IP address and port combination.
Definition IPEndpoint.h:24
Endpoint atPort(Port port) const
Returns a new Endpoint with a different port.
Definition IPEndpoint.h:65
Address const & address() const
Returns the address portion of this endpoint.
Definition IPEndpoint.h:74
Represents a JSON value.
Definition json_value.h:117
A consumption charge.
Definition Charge.h:13
int value_type
The type used to hold a consumption charge.
Definition Charge.h:18
value_type cost() const
Return the cost of the charge in resource::Manager units.
Definition Charge.cpp:22
An endpoint that consumes resources.
Definition Consumer.h:20
static void writeList(ClockType::time_point const now, beast::PropertyStream::Set &items, EntryIntrusiveList &list)
HashMap< std::string, Import > Imports
HashMap< Key, Entry, Key::Hasher, Key::KeyEqual > Table
bool disconnect(Entry &entry)
Disposition charge(Entry &entry, Charge const &fee, std::string context={})
json::Value getJson(int threshold)
Returns a json::ValueType::Object.
void erase(Table::iterator iter)
Consumer newInboundEndpoint(beast::ip::Endpoint const &address)
void onWrite(beast::PropertyStream::Map &map)
static Disposition disposition(int balance)
Consumer newOutboundEndpoint(beast::ip::Endpoint const &address)
Logic(beast::insight::Collector::Ptr const &collector, ClockType &clock, beast::Journal journal)
EntryIntrusiveList inactive_
std::recursive_mutex lock_
EntryIntrusiveList inbound_
Consumer newUnlimitedEndpoint(beast::ip::Endpoint const &address)
Create endpoint that should not have resource limits applied.
beast::List< Entry > EntryIntrusiveList
void importConsumers(std::string const &origin, Gossip const &gossip)
EntryIntrusiveList outbound_
T make_tuple(T... args)
@ Object
object value (collection of name/value pairs).
Definition json_value.h:29
static constexpr auto kMinimumGossipBalance
static constexpr std::chrono::seconds kSecondsUntilExpiration
static constexpr auto kWarningThreshold
Tunable constants.
Disposition
The disposition of a consumer after applying a load charge.
Definition Disposition.h:8
@ Warn
Consumer should be disconnected for excess consumption.
Definition Disposition.h:18
@ Ok
No action required.
Definition Disposition.h:12
Charge const kFeeDrop
Charge const kFeeWarning
static constexpr auto kDropThreshold
static constexpr std::chrono::seconds kGossipExpirationSeconds
beast::AbstractClock< std::chrono::steady_clock > Stopwatch
A clock for measuring elapsed time.
Definition chrono.h:90
std::unordered_map< Key, Value, Hash, Pred, Allocator > HashMap
T piecewise_construct
Describes a single consumer.
Definition Gossip.h:20
beast::ip::Endpoint address
Definition Gossip.h:24
Data format for exchanging consumption information across peers.
Definition Gossip.h:13
std::vector< Item > items
Definition Gossip.h:27
A set of imported consumer data from a gossip origin.
Definition Import.h:14
Stats(beast::insight::Collector::Ptr const &collector)
T swap(T... args)