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>
41 warn = collector->makeMeter(
"warn");
42 drop = collector->makeMeter(
"drop");
98 Entry* entry(
nullptr);
102 auto [resultIt, resultInserted] =
table_.emplace(
107 entry = &resultIt->second;
108 entry->key = &resultIt->first;
110 if (entry->refcount == 1)
120 JLOG(
journal_.debug()) <<
"New inbound endpoint " << *entry;
128 Entry* entry(
nullptr);
132 auto [resultIt, resultInserted] =
table_.emplace(
137 entry = &resultIt->second;
138 entry->key = &resultIt->first;
140 if (entry->refcount == 1)
148 JLOG(
journal_.debug()) <<
"New outbound endpoint " << *entry;
161 Entry* entry(
nullptr);
165 auto [resultIt, resultInserted] =
table_.emplace(
170 entry = &resultIt->second;
171 entry->key = &resultIt->first;
173 if (entry->refcount == 1)
181 JLOG(
journal_.debug()) <<
"New unlimited endpoint " << *entry;
205 int const localBalance = inboundEntry.localBalance.value(now);
206 if ((localBalance + inboundEntry.remoteBalance) >= threshold)
209 entry[jss::local] = localBalance;
210 entry[jss::remote] = inboundEntry.remoteBalance;
211 entry[jss::type] =
"inbound";
216 int const localBalance = outboundEntry.localBalance.value(now);
217 if ((localBalance + outboundEntry.remoteBalance) >= threshold)
220 entry[jss::local] = localBalance;
221 entry[jss::remote] = outboundEntry.remoteBalance;
222 entry[jss::type] =
"outbound";
225 for (
auto& adminEntry :
admin_)
227 int const localBalance = adminEntry.localBalance.value(now);
228 if ((localBalance + adminEntry.remoteBalance) >= threshold)
231 entry[jss::local] = localBalance;
232 entry[jss::remote] = adminEntry.remoteBalance;
233 entry[jss::type] =
"admin";
253 item.
balance = inboundEntry.localBalance.value(now);
257 gossip.
items.push_back(item);
269 auto const elapsed =
clock_.now();
280 Import& next(resultIt->second);
282 next.items.reserve(gossip.
items.size());
284 for (
auto const& gossipItem : gossip.
items)
287 item.
balance = gossipItem.balance;
290 next.items.push_back(item);
300 next.items.reserve(gossip.
items.size());
301 for (
auto const& gossipItem : gossip.
items)
304 item.
balance = gossipItem.balance;
307 next.items.push_back(item);
310 Import& prev(resultIt->second);
311 for (
auto& item : prev.items)
313 item.consumer.entry().remoteBalance -= item.balance;
330 auto const elapsed =
clock_.now();
334 if (iter->whenExpires <= elapsed)
336 JLOG(
journal_.debug()) <<
"Expired " << *iter;
337 auto tableIter =
table_.find(*iter->key);
350 Import&
import(iter->second);
351 if (iter->second.whenExpires <= elapsed)
353 for (
auto& item :
import.items)
355 item.consumer.entry().remoteBalance -= item.balance;
386 Entry& entry(iter->second);
387 XRPL_ASSERT(entry.refcount == 0,
"xrpl::resource::Logic::erase : entry not used");
403 if (--entry.refcount == 0)
405 JLOG(
journal_.debug()) <<
"Inactive " << entry;
407 switch (entry.key->kind)
421 "xrpl::resource::Logic::release : invalid entry "
438 kFeeLogAsWarn > kFeeLogAsInfo && kFeeLogAsInfo > kFeeLogAsDebug && kFeeLogAsDebug > 10);
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();
450 if (!context.empty())
451 context =
" (" + context +
")";
456 JLOG(kGetStream(fee.
cost(),
journal_)) <<
"Charging " << entry <<
" for " << fee << context;
463 if (entry.isUnlimited())
468 auto const elapsed =
clock_.now();
473 entry.lastWarningTime = elapsed;
477 JLOG(
journal_.info()) <<
"Load warning: " << entry;
486 if (entry.isUnlimited())
492 int const balance(entry.balance(now));
495 JLOG(
journal_.warn()) <<
"Consumer entry " << entry <<
" dropped with balance "
512 return entry.balance(
clock_.now());
523 for (
auto& entry : list)
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;
std::chrono::steady_clock::time_point time_point
virtual time_point now() const =0
Returns the current time.
A generic endpoint for log messages.
Intrusive doubly linked list.
std::shared_ptr< Collector > Ptr
A metric for measuring an integral value.
A version-independent IP address and port combination.
Endpoint atPort(Port port) const
Returns a new Endpoint with a different port.
Address const & address() const
Returns the address portion of this endpoint.
int value_type
The type used to hold a consumption charge.
value_type cost() const
Return the cost of the charge in resource::Manager units.
An endpoint that consumes resources.
static void writeList(ClockType::time_point const now, beast::PropertyStream::Set &items, EntryIntrusiveList &list)
HashMap< std::string, Import > Imports
void release(Entry &entry)
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)
void acquire(Entry &entry)
Logic(beast::insight::Collector::Ptr const &collector, ClockType &clock, beast::Journal journal)
EntryIntrusiveList inactive_
std::recursive_mutex lock_
EntryIntrusiveList inbound_
EntryIntrusiveList admin_
Consumer newUnlimitedEndpoint(beast::ip::Endpoint const &address)
Create endpoint that should not have resource limits applied.
int balance(Entry &entry)
beast::List< Entry > EntryIntrusiveList
void importConsumers(std::string const &origin, Gossip const &gossip)
EntryIntrusiveList outbound_
@ Object
object value (collection of name/value pairs).
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.
@ Warn
Consumer should be disconnected for excess consumption.
static constexpr auto kDropThreshold
static constexpr std::chrono::seconds kGossipExpirationSeconds
beast::AbstractClock< std::chrono::steady_clock > Stopwatch
A clock for measuring elapsed time.
std::unordered_map< Key, Value, Hash, Pred, Allocator > HashMap
Describes a single consumer.
beast::ip::Endpoint address
Data format for exchanging consumption information across peers.
std::vector< Item > items
A set of imported consumer data from a gossip origin.
beast::insight::Meter drop
Stats(beast::insight::Collector::Ptr const &collector)
beast::insight::Meter warn