Merge branch 'pratik/otel-phase10-workload-validation' into pratik/otel-sync-diagnostics

Brings the MetricsRegistry split onto this branch. The pipeline half is now
xrpl::telemetry::MetricsRegistry in libxrpl; the observable gauges are
xrpl::telemetry::AppMetricGauges in xrpld.

This branch had added its own instrumentation to the pre-split class, so the
merge had to route each addition to the correct half:

- The thirteen gauges added here -- amendment block, cache hit-rate detail,
  clock skew, job-queue saturation, ledger quorum publish, peer ledger supply,
  rotation state, slot census, stall events, sync acquire, sync state, UNL
  quorum, and the cache lock-hold observer -- all land on AppMetricGauges,
  reading the core's meter and validation tracker through it.
- The pipeline additions stay in libxrpl: the consensus round-duration and
  rotation-phase histogram views, the malloc-trim and dns/dial latency bucket
  ladders, the job-stall counter, and the switch from literal metric names to
  the MetricNames.h constants.

Git detected the pre-split MetricsRegistry.cpp and .h as renames of the gauge
files, so both sides' pipeline changes initially landed in the gauge half. They
were moved back, and the result was audited by inventory: every method
definition, instrument creation, view registration, and emitted string from
either side is present, with identical multiplicity.

MetricNames.h moves to include/xrpl/telemetry/ alongside the core. It has no
includes of its own and its two sibling name headers already live there, so
keeping it under src/ would leave an xrpld path in libxrpl's dependency
surface. Nineteen files follow it.

incrementStateChanges() stays removed. The labelled state_changes_total{from,to}
counter this branch introduced replaces it, and the test asserting the method is
absent is kept -- an unlabelled instrument alongside the labelled one would give
Prometheus two conflicting versions of one metric name.

Two tests that drove startAsyncGauges() against a mock ServiceRegistry are
dropped: xrpl_tests links only xrpl.libxrpl and cannot reach the gauge class.

Levelization regenerated. Both xrpld.telemetry loops become bidirectional
rather than one-way; neither is new.
This commit is contained in:
Pratik Mankawde
2026-09-16 18:05:48 +01:00
57 changed files with 3637 additions and 3592 deletions

View File

@@ -0,0 +1,449 @@
#pragma once
/**
* Call-site OTel metric macros.
*
* Adds a new OTel metric instrument entirely at the call site -- no member
* field, no init line, no wrapper method in MetricsRegistry. Covers every
* instrument kind the OTel Metrics API defines:
*
* Synchronous (created once on first use, then record on every call):
* Counter XRPL_METRIC_COUNTER_INC / _ADD [+ _LABELED]
* UpDownCounter XRPL_METRIC_UPDOWN_ADD [+ _LABELED]
* Histogram XRPL_METRIC_HISTOGRAM_RECORD [+ _LABELED]
* Gauge XRPL_METRIC_GAUGE_RECORD [+ _LABELED]
* (requires OPENTELEMETRY_ABI_VERSION_NO >= 2;
* this repo currently builds ABI v1 -- see the
* static_assert branch below)
*
* Asynchronous/observable (register a callback ONCE, eagerly, during
* construction/init -- see the "Observable" note below):
* Observable Counter XRPL_METRIC_OBSERVABLE_COUNTER_REGISTER
* Observable UpDownCounter XRPL_METRIC_OBSERVABLE_UPDOWN_REGISTER
* Observable Gauge XRPL_METRIC_OBSERVABLE_GAUGE_REGISTER
*
* When XRPL_ENABLE_TELEMETRY is not defined, every macro expands to a
* no-op statement, so call sites never need their own #ifdef.
*
* Example usage -- plain counter:
* @code
* void RCLConsensus::Adaptor::doAccept()
* {
* // ... existing consensus-accept logic ...
* XRPL_METRIC_COUNTER_INC(app_, "ledgers_closed_total",
* "Total ledgers closed by consensus");
* }
* @endcode
*
* Example usage -- labeled counter (edge case: per-reason tally):
* @code
* void TxQ::apply(...)
* {
* if (queueIsFull)
* XRPL_METRIC_COUNTER_INC_LABELED(app, "txq_dropped_total",
* "Transactions refused admission to the queue",
* {{"reason", std::string("queue_full")}});
* }
* @endcode
*
* Pass the label set as a bare brace-enclosed list, as above. Do not wrap
* it in an extra pair of parentheses: the list is forwarded verbatim into
* the OTel `Add()`/`Record()` call, which takes an initializer_list, and
* the extra parentheses do not compile.
*
* Wrap each label *value* in `std::string`. `AttributeValue` is a variant
* in which a bare `const char*` selects the boolean alternative, so an
* unwrapped literal is recorded as `true`.
*
* Example usage -- UpDownCounter (edge case: value that can decrease):
* @code
* void ServerHandler::onRpcStart()
* {
* XRPL_METRIC_UPDOWN_ADD(app_, "rpc_in_flight_requests",
* "RPC requests currently executing", 1);
* }
* void ServerHandler::onRpcFinish()
* {
* XRPL_METRIC_UPDOWN_ADD(app_, "rpc_in_flight_requests",
* "RPC requests currently executing", -1);
* }
* @endcode
*
* Example usage -- observable gauge registered from a non-MetricsRegistry
* class (edge case: a subsystem exposing its own live state):
* @code
* SomeSubsystem::SomeSubsystem(ServiceRegistry& app) : app_(app)
* {
* XRPL_METRIC_OBSERVABLE_GAUGE_REGISTER(
* app_, "some_subsystem_queue_depth", "Current queue depth",
* [this] { return static_cast<int64_t>(queue_.size()); });
* }
* @endcode
*
* @note A histogram whose values can exceed ~10,000 units (e.g. a
* microsecond duration beyond 10ms) needs an explicit-bucket View, which
* OTel can only register at MeterProvider construction time -- this
* cannot be done from a call site. Register such a view in
* MetricsRegistry::initExporterAndProvider() as today; the
* histogram-record call itself can still use the macro.
*
* @note The SYNCHRONOUS macros (Counter/UpDownCounter/Histogram/Gauge)
* create their instrument once, on first use, from
* MetricsRegistry::meter(). The registry builds that meter in its
* constructor, before any subsystem exists, and guarantees it is never
* empty while the registry is enabled (a no-op meter stands in if the
* pipeline failed to build, and again after stop()). So a call site holds
* a valid instrument from its first call and needs no check of its own.
* The only branch on the hot path is the recording() gate, which is false
* once stop() has torn the pipeline down; without that gate a Record on a
* stale SDK instrument would deref a dangling AggregationConfig.
*
* @note Static-init safety: Meter::CreateXxx is declared noexcept in the
* OTel API (opentelemetry/metrics/meter.h), so the function-local static
* that caches the instrument cannot throw during first-call construction.
* A throw there would call std::terminate.
*
* @note The OBSERVABLE registration macros are the opposite: call them
* EAGERLY, exactly once, from constructor/init code -- never from a hot
* path. Repeated calls at the same call site register a NEW callback
* each time (no create-once caching, unlike the synchronous macros),
* which leaks callbacks.
*
* @note There is no way to read back a synchronous instrument's current
* accumulated value from application code -- the OTel API is
* write-only/push-based by design. If your logic needs both to record a
* metric AND read its running value, keep your own state (std::atomic or
* similar) and separately feed OTel via these macros.
*/
// On Windows, OTel's spin_lock_mutex.h (transitively included from
// MetricsRegistry.h) defines _WINSOCKAPI_ and includes <windows.h>, which
// pulls in WinSock 1. ServiceRegistry.h below then includes <boost/asio.hpp>,
// whose socket_types.hpp requires winsock2.h first and errors out if WinSock 1
// arrived earlier. Pre-including boost's socket types header here gets
// winsock2.h in before the OTel headers, so any translation unit that includes
// MetricMacros.h first (e.g. the telemetry unit tests) still compiles. The
// production MetricsRegistry.cpp carries the same guard.
#ifdef _MSC_VER
#include <boost/asio/detail/socket_types.hpp>
#endif
#include <xrpl/core/ServiceRegistry.h> // IWYU pragma: keep
#include <xrpl/telemetry/MetricsRegistry.h> // IWYU pragma: keep
#ifdef XRPL_ENABLE_TELEMETRY
#include <functional> // IWYU pragma: keep
#define XRPL_METRIC_COUNTER_INC(app, name, description) \
do \
{ \
if (auto* xrpl_mr_ = (app).getMetricsRegistry(); xrpl_mr_ && xrpl_mr_->recording()) \
{ \
static auto const xrpl_counter_ = \
xrpl_mr_->meter()->CreateUInt64Counter((name), (description)); \
xrpl_counter_->Add(1); \
} \
} while (false)
// The label set is passed as trailing variadic arguments so a
// brace-enclosed initializer list (e.g. {{"reason", std::string("x")}}),
// which contains a top-level comma, survives preprocessing as a single
// logical argument. __VA_ARGS__ re-joins it verbatim into the Add() call.
#define XRPL_METRIC_COUNTER_INC_LABELED(app, name, description, ...) \
do \
{ \
if (auto* xrpl_mr_ = (app).getMetricsRegistry(); xrpl_mr_ && xrpl_mr_->recording()) \
{ \
static auto const xrpl_counter_ = \
xrpl_mr_->meter()->CreateUInt64Counter((name), (description)); \
xrpl_counter_->Add(1, __VA_ARGS__); \
} \
} while (false)
// Same as XRPL_METRIC_COUNTER_INC, but increments by a caller-supplied amount
// instead of a fixed 1 (e.g. bytes transferred, batch sizes).
#define XRPL_METRIC_COUNTER_ADD(app, name, description, amount) \
do \
{ \
if (auto* xrpl_mr_ = (app).getMetricsRegistry(); xrpl_mr_ && xrpl_mr_->recording()) \
{ \
static auto const xrpl_counter_ = \
xrpl_mr_->meter()->CreateUInt64Counter((name), (description)); \
xrpl_counter_->Add(amount); \
} \
} while (false)
// amount is fixed; the trailing variadic args carry the label set (see the
// note on XRPL_METRIC_COUNTER_INC_LABELED for why labels are variadic).
#define XRPL_METRIC_COUNTER_ADD_LABELED(app, name, description, amount, ...) \
do \
{ \
if (auto* xrpl_mr_ = (app).getMetricsRegistry(); xrpl_mr_ && xrpl_mr_->recording()) \
{ \
static auto const xrpl_counter_ = \
xrpl_mr_->meter()->CreateUInt64Counter((name), (description)); \
xrpl_counter_->Add(amount, __VA_ARGS__); \
} \
} while (false)
// UpDownCounter: like COUNTER_ADD, but the underlying instrument permits a
// negative amount (e.g. in-flight request count, +1 on start / -1 on
// finish from two different points in the same or different call sites).
// A plain Counter's Add() must never see a negative value per the OTel
// API contract; use this macro, not COUNTER_ADD, whenever the value can
// decrease.
#define XRPL_METRIC_UPDOWN_ADD(app, name, description, amount) \
do \
{ \
if (auto* xrpl_mr_ = (app).getMetricsRegistry(); xrpl_mr_ && xrpl_mr_->recording()) \
{ \
static auto const xrpl_updown_ = \
xrpl_mr_->meter()->CreateInt64UpDownCounter((name), (description)); \
xrpl_updown_->Add(amount); \
} \
} while (false)
// amount may be negative; the trailing variadic args carry the label set
// (see the note on XRPL_METRIC_COUNTER_INC_LABELED for why labels are variadic).
#define XRPL_METRIC_UPDOWN_ADD_LABELED(app, name, description, amount, ...) \
do \
{ \
if (auto* xrpl_mr_ = (app).getMetricsRegistry(); xrpl_mr_ && xrpl_mr_->recording()) \
{ \
static auto const xrpl_updown_ = \
xrpl_mr_->meter()->CreateInt64UpDownCounter((name), (description)); \
xrpl_updown_->Add(amount, __VA_ARGS__); \
} \
} while (false)
#define XRPL_METRIC_HISTOGRAM_RECORD(app, name, description, value) \
do \
{ \
if (auto* xrpl_mr_ = (app).getMetricsRegistry(); xrpl_mr_ && xrpl_mr_->recording()) \
{ \
static auto const xrpl_hist_ = \
xrpl_mr_->meter()->CreateDoubleHistogram((name), (description)); \
xrpl_hist_->Record(static_cast<double>(value), opentelemetry::context::Context{}); \
} \
} while (false)
// value is fixed; the trailing variadic args carry the label set (see the
// note on XRPL_METRIC_COUNTER_INC_LABELED for why labels are variadic).
#define XRPL_METRIC_HISTOGRAM_RECORD_LABELED(app, name, description, value, ...) \
do \
{ \
if (auto* xrpl_mr_ = (app).getMetricsRegistry(); xrpl_mr_ && xrpl_mr_->recording()) \
{ \
static auto const xrpl_hist_ = \
xrpl_mr_->meter()->CreateDoubleHistogram((name), (description)); \
xrpl_hist_->Record( \
static_cast<double>(value), __VA_ARGS__, opentelemetry::context::Context{}); \
} \
} while (false)
// Synchronous Gauge: last-value snapshot, not a distribution (contrast
// Histogram) and not a running total (contrast Counter/UpDownCounter).
// ABI-gated: opentelemetry-cpp only exposes CreateInt64Gauge/
// CreateDoubleGauge when OPENTELEMETRY_ABI_VERSION_NO >= 2. This project's
// Conan build currently pins ABI v1 (verified:
// .build/build/generators/opentelemetry-cpp-release-x86_64-data.cmake sets
// OPENTELEMETRY_ABI_VERSION_NO=1), so the real path below is presently
// dead code on this codebase's build -- shipped anyway so it activates
// automatically the day the ABI version is bumped, and so a developer who
// reaches for "just the current value, not a distribution" sees an
// actionable compile error now instead of silently reaching for the wrong
// instrument kind (misusing Histogram or UpDownCounter as a gauge
// substitute is explicitly discouraged -- see Design/taxonomy section).
#if OPENTELEMETRY_ABI_VERSION_NO >= 2
#define XRPL_METRIC_GAUGE_RECORD(app, name, description, value) \
do \
{ \
if (auto* xrpl_mr_ = (app).getMetricsRegistry(); xrpl_mr_ && xrpl_mr_->recording()) \
{ \
static auto const xrpl_gauge_ = \
xrpl_mr_->meter()->CreateDoubleGauge((name), (description)); \
xrpl_gauge_->Record(static_cast<double>(value), opentelemetry::context::Context{}); \
} \
} while (false)
// value is fixed; the trailing variadic args carry the label set (see the
// note on XRPL_METRIC_COUNTER_INC_LABELED for why labels are variadic).
#define XRPL_METRIC_GAUGE_RECORD_LABELED(app, name, description, value, ...) \
do \
{ \
if (auto* xrpl_mr_ = (app).getMetricsRegistry(); xrpl_mr_ && xrpl_mr_->recording()) \
{ \
static auto const xrpl_gauge_ = \
xrpl_mr_->meter()->CreateDoubleGauge((name), (description)); \
xrpl_gauge_->Record( \
static_cast<double>(value), __VA_ARGS__, opentelemetry::context::Context{}); \
} \
} while (false)
#else
#define XRPL_METRIC_GAUGE_RECORD(app, name, description, value) \
static_assert( \
false, \
"XRPL_METRIC_GAUGE_RECORD requires OPENTELEMETRY_ABI_VERSION_NO >= 2 (this build " \
"uses ABI v1). Use XRPL_METRIC_OBSERVABLE_GAUGE_REGISTER with your own state, or " \
"bump OPENTELEMETRY_ABI_VERSION_NO as a separate, reviewed change.")
#define XRPL_METRIC_GAUGE_RECORD_LABELED(app, name, description, value, ...) \
static_assert( \
false, \
"XRPL_METRIC_GAUGE_RECORD_LABELED requires OPENTELEMETRY_ABI_VERSION_NO >= 2 " \
"(this build uses ABI v1). Use XRPL_METRIC_OBSERVABLE_GAUGE_REGISTER with your " \
"own state, or bump OPENTELEMETRY_ABI_VERSION_NO as a separate, reviewed change.")
#endif // OPENTELEMETRY_ABI_VERSION_NO >= 2
// -----------------------------------------------------------------
// Observable/async instrument registration. Unlike the synchronous
// macros above, these do NOT lazily create-on-first-call -- they
// register a callback with the SDK immediately, at the call site, the
// moment the macro executes. Callers MUST invoke this during
// construction/init, before the server is fully live (same timing rule
// MetricsRegistry::registerAsyncGauges() already follows for its own
// gauges). Calling it from a hot-path function instead of an init path
// re-registers a new callback on every call, which leaks callbacks and
// is NOT what this macro is for.
//
// The callable is captured in a heap-allocated std::function, and its
// address is passed as the `void* state` to AddCallback (whose signature,
// `void (*)(ObserverResult, void*)`, is a raw C function pointer -- it
// cannot bind a capturing lambda directly). A static trampoline
// function reinterprets `state` back to the std::function and invokes
// it inside the callback. The heap allocation is intentionally leaked
// for the process lifetime (matches every existing ObservableGauge in
// MetricsRegistry, which are member fields with the same lifetime as
// the registry itself) -- do not "fix" this with a smart pointer that
// frees before the reader thread's last collection tick.
// -----------------------------------------------------------------
#define XRPL_METRIC_OBSERVABLE_GAUGE_REGISTER(app, name, description, valueFn) \
do \
{ \
if (auto* xrpl_mr_ = (app).getMetricsRegistry(); xrpl_mr_ && xrpl_mr_->recording()) \
{ \
auto xrpl_m_ = xrpl_mr_->meter(); \
auto* xrpl_fn_ = new std::function<int64_t()>(valueFn); \
auto xrpl_inst_ = xrpl_m_->CreateInt64ObservableGauge((name), (description)); \
xrpl_inst_->AddCallback( \
[](opentelemetry::metrics::ObserverResult result, void* state) { \
auto* fn = static_cast<std::function<int64_t()>*>(state); \
try \
{ \
opentelemetry::nostd::get<opentelemetry::nostd::shared_ptr< \
opentelemetry::metrics::ObserverResultT<int64_t>>>(result) \
->Observe((*fn)()); \
} \
catch (...) \
{ \
} \
}, \
xrpl_fn_); \
} \
} while (false)
#define XRPL_METRIC_OBSERVABLE_COUNTER_REGISTER(app, name, description, valueFn) \
do \
{ \
if (auto* xrpl_mr_ = (app).getMetricsRegistry(); xrpl_mr_ && xrpl_mr_->recording()) \
{ \
auto xrpl_m_ = xrpl_mr_->meter(); \
auto* xrpl_fn_ = new std::function<int64_t()>(valueFn); \
auto xrpl_inst_ = xrpl_m_->CreateInt64ObservableCounter((name), (description)); \
xrpl_inst_->AddCallback( \
[](opentelemetry::metrics::ObserverResult result, void* state) { \
auto* fn = static_cast<std::function<int64_t()>*>(state); \
try \
{ \
opentelemetry::nostd::get<opentelemetry::nostd::shared_ptr< \
opentelemetry::metrics::ObserverResultT<int64_t>>>(result) \
->Observe((*fn)()); \
} \
catch (...) \
{ \
} \
}, \
xrpl_fn_); \
} \
} while (false)
#define XRPL_METRIC_OBSERVABLE_UPDOWN_REGISTER(app, name, description, valueFn) \
do \
{ \
if (auto* xrpl_mr_ = (app).getMetricsRegistry(); xrpl_mr_ && xrpl_mr_->recording()) \
{ \
auto xrpl_m_ = xrpl_mr_->meter(); \
auto* xrpl_fn_ = new std::function<int64_t()>(valueFn); \
auto xrpl_inst_ = xrpl_m_->CreateInt64ObservableUpDownCounter((name), (description)); \
xrpl_inst_->AddCallback( \
[](opentelemetry::metrics::ObserverResult result, void* state) { \
auto* fn = static_cast<std::function<int64_t()>*>(state); \
try \
{ \
opentelemetry::nostd::get<opentelemetry::nostd::shared_ptr< \
opentelemetry::metrics::ObserverResultT<int64_t>>>(result) \
->Observe((*fn)()); \
} \
catch (...) \
{ \
} \
}, \
xrpl_fn_); \
} \
} while (false)
#else // !XRPL_ENABLE_TELEMETRY
#define XRPL_METRIC_COUNTER_INC(app, name, description) \
do \
{ \
} while (false)
#define XRPL_METRIC_COUNTER_INC_LABELED(app, name, description, ...) \
do \
{ \
} while (false)
#define XRPL_METRIC_COUNTER_ADD(app, name, description, amount) \
do \
{ \
} while (false)
#define XRPL_METRIC_COUNTER_ADD_LABELED(app, name, description, amount, ...) \
do \
{ \
} while (false)
#define XRPL_METRIC_UPDOWN_ADD(app, name, description, amount) \
do \
{ \
} while (false)
#define XRPL_METRIC_UPDOWN_ADD_LABELED(app, name, description, amount, ...) \
do \
{ \
} while (false)
#define XRPL_METRIC_HISTOGRAM_RECORD(app, name, description, value) \
do \
{ \
} while (false)
#define XRPL_METRIC_HISTOGRAM_RECORD_LABELED(app, name, description, value, ...) \
do \
{ \
} while (false)
#define XRPL_METRIC_GAUGE_RECORD(app, name, description, value) \
do \
{ \
} while (false)
#define XRPL_METRIC_GAUGE_RECORD_LABELED(app, name, description, value, ...) \
do \
{ \
} while (false)
#define XRPL_METRIC_OBSERVABLE_GAUGE_REGISTER(app, name, description, valueFn) \
do \
{ \
} while (false)
#define XRPL_METRIC_OBSERVABLE_COUNTER_REGISTER(app, name, description, valueFn) \
do \
{ \
} while (false)
#define XRPL_METRIC_OBSERVABLE_UPDOWN_REGISTER(app, name, description, valueFn) \
do \
{ \
} while (false)
#endif // XRPL_ENABLE_TELEMETRY

View File

@@ -0,0 +1,794 @@
#pragma once
/**
* Compile-time OTel metric name constants for the sync-diagnostics signals.
*
* The metric-side counterpart of the `*SpanNames.h` headers: one constant per
* emitted string, so a rename is one edit and a typo is a compile error rather
* than a metric that silently never appears. Each instrument name, label key
* and bounded label value added by the sync-diagnostics work is declared here
* exactly once and referenced from every C++ user of it -- the emit site, the
* gauge registration in MetricsRegistry.cpp, and the unit test.
*
* Layer map -- who references these constants:
*
* +-------------------------------------------------------------+
* | MetricNames.h (this file, L1-metrics) |
* | namespace metric namespace label namespace lval |
* +-------------------------------------------------------------+
* ^ ^ ^
* | | |
* +--------------+ +-------------------+ +---------------------+
* | emit sites | | MetricsRegistry | | unit test |
* | XRPL_METRIC_ | | Create*Gauge + | | tests/libxrpl/ |
* | macros under | | AddCallback | | telemetry/ |
* | src/xrpld/ | | observe(...) | | MetricMacros.cpp |
* +--------------+ +-------------------+ +---------------------+
*
* Layers that CANNOT reference a C++ constant -- the collector config, the
* dashboard PromQL, `expected_metrics.json` and the runbook -- are held to
* these same strings by `.github/scripts/otel-naming/check_otel_naming.py`
* instead.
*
* Where this is documented:
* - CONTRIBUTING.md -> "Telemetry metric naming" is the authoritative rule
* list, alongside the sibling span-attribute convention.
* - `.github/scripts/otel-naming/README.md` describes the enforcing rules
* I (no literals), J (suffix conventions) and K (expected_metrics.json).
* - docs/telemetry-runbook.md -> "Adding a New Metric" is the walkthrough for
* adding one, and shows the declare-then-emit pattern.
*
* Naming rules (enforced by the checker's Rule J):
* - Bare `lower_snake_case`. No `xrpld_` prefix in code: the Prometheus
* exporter adds the namespace prefix itself, so writing it here would
* produce `xrpld_xrpld_*` on the wire.
* - A monotonic counter ends in `_total`, so `rate()` over it reads correctly
* and a reader can tell it from a gauge at a glance.
* - A duration carries its unit as the suffix: `_us`, `_ms` or `_seconds`.
* The unit belongs in the name because the OTel `unit` argument is not
* surfaced on the Prometheus metric name.
* - A gauge that is a snapshot of current state takes no suffix
* (`jobq_saturation`, `sync_state`).
* - Label VALUES are declared here only when they come from a fixed set that
* the code itself writes (`namespace lval`), which is what keeps series
* cardinality bounded. A value derived from runtime data -- a site URI, a
* peer address, a ledger hash -- is deliberately NOT declared here and must
* never become a label on a metric.
*
* Why `constexpr char[]` and not the `makeStr`/`StaticStr` DSL that
* `SpanNames.h` uses: the OTel C++ API takes `nostd::string_view`, which in
* this build (`OPENTELEMETRY_STL_VERSION` unset, so the back-ported class is
* used) constructs only from `char const*`, `std::string` or
* `(char const*, size)`. It has NO constructor from `std::string_view`, so
* neither `StaticStr` (which converts to `std::string_view`) nor a
* `constexpr std::string_view` compiles in an instrument-name or label-key
* position -- both were tried and both fail with "no viable conversion".
* A `constexpr char[]` decays to `char const*` and binds directly. This also
* matches the precedent already in MetricsRegistry.cpp
* (`kJobQueuedDurationUs`), which this header absorbs.
*
* Example usage -- a labelled counter at an emit site:
* @code
* XRPL_METRIC_COUNTER_INC_LABELED(
* app_,
* metric::dnsResolveTotal,
* "Peer hostname resolutions, by outcome",
* {{label::outcome,
* std::string(
* resolved ? lval::dns_resolve::resolved : lval::dns_resolve::empty)}});
* @endcode
*
* Example usage -- an observable gauge and its sub-metric discriminators:
* @code
* syncStateGauge_ = meter_->CreateInt64ObservableGauge(
* metric::syncState, "Sync-pipeline health signals");
* // ... inside the callback:
* observe(lval::sync_state::ledgersBehind, ops.getLedgersBehindNetwork());
* @endcode
*
* Example usage -- edge case: a value that must NOT be a constant. The site
* URI is runtime data, so only the KEY is named here; declaring the value
* would imply a bounded set that does not exist. Note the value is built from
* the PARSED url, not the configured string: the config accepts userinfo and a
* query, and either would otherwise ride the label into Prometheus.
* @code
* auto const& url = sites_[siteIdx].loadedResource->pUrl;
* std::string site = url.scheme + "://" + url.domain;
* site += url.path.substr(0, url.path.find_first_of("?#"));
* XRPL_METRIC_COUNTER_INC_LABELED(
* app_, metric::unlFetchTotal, "...",
* {{label::site, site}, {label::outcome, std::string(outcome)}});
* @endcode
*
* @note Header-only and dependency-free: nothing here includes an OTel or an
* xrpld header, so `src/tests/libxrpl/telemetry/MetricMacros.cpp` can
* include it even though `xrpl_tests` links only `xrpl.libxrpl`. The
* constants are `inline constexpr`, so they contribute no symbol to link
* against.
* @note Not guarded by `XRPL_ENABLE_TELEMETRY`, for the same reason
* `SpanNames.h` is not: call sites name these constants even where the
* macros expand to no-ops, and the compiler elides any constant whose
* only uses are in dead code.
* @note Every constant is a compile-time value with no mutable state, so all
* of them are safe to read concurrently from any thread.
*/
namespace xrpl::telemetry {
/**
* Instrument names -- the metric name as it reaches the OTel meter.
*
* Grouped by the subsystem that emits them, matching how the sync-diagnostics
* work was staged. A name is declared here whether it is created lazily by an
* `XRPL_METRIC_*` macro at a call site or eagerly by a `meter_->Create*` call
* in MetricsRegistry.cpp, because the dashboards cannot tell the two apart.
*/
namespace metric {
// ===== Bootstrap: getting a fresh node its first peers and its first UNL =====
/**
* Time to resolve one configured peer hostname.
*/
inline constexpr char dnsResolveLatencyMs[] = "dns_resolve_latency_ms";
/**
* Peer hostname resolutions, split by whether any address came back.
*/
inline constexpr char dnsResolveTotal[] = "dns_resolve_total";
/**
* Time from starting an outbound dial to its terminal outcome.
*/
inline constexpr char overlayDialLatencyMs[] = "overlay_dial_latency_ms";
/**
* Outbound peer connection attempts, by terminal outcome.
*/
inline constexpr char overlayConnectTotal[] = "overlay_connect_total";
/**
* Peer handshakes this node rejected, by reason.
*/
inline constexpr char handshakeNegotiationFailTotal[] = "handshake_negotiation_fail_total";
/**
* Validator-list fetch attempts, by site and outcome.
*/
inline constexpr char unlFetchTotal[] = "unl_fetch_total";
/**
* Trusted UNL key count against the required quorum.
*
* A gauge, not a counter: both series are current state, and the useful read
* is the difference between them.
*/
inline constexpr char unlQuorum[] = "unl_quorum";
/**
* Network close-time offset from the local clock.
*/
inline constexpr char clockCloseOffsetSeconds[] = "clock_close_offset_seconds";
// ===== Sync state: why this node is not FULL yet =============================
/**
* Operating-mode transitions, labelled with the mode pair.
*/
inline constexpr char stateChangesTotal[] = "state_changes_total";
/**
* Four sync-pipeline health signals, split by the `metric` label.
*/
inline constexpr char syncState[] = "sync_state";
/**
* Main-loop stall episodes. Cumulative, so `rate()` gives episodes/sec.
*/
inline constexpr char serverStallEventsTotal[] = "server_stall_events_total";
// ===== Acquire / SHAMap: is ledger data actually arriving? ===================
/**
* Aggregate ledger-acquire progress across all in-flight acquires.
*/
inline constexpr char syncAcquire[] = "sync_acquire";
/**
* SHAMap tree-node cache hit rate, 0.0-1.0.
*/
inline constexpr char shamapCacheHitRate[] = "shamap_cache_hit_rate";
/**
* Ledger acquires by where the data came from (local store vs network).
*/
inline constexpr char syncAcquireSourceTotal[] = "sync_acquire_source_total";
/**
* Acquire timeouts where not one new node arrived.
*/
inline constexpr char syncAcquireNoProgressTotal[] = "sync_acquire_no_progress_total";
/**
* SHAMap nodes received during an acquire, by per-node outcome.
*/
inline constexpr char syncAddnodeTotal[] = "sync_addnode_total";
// ===== JobQueue: is the worker pool the bottleneck? ==========================
/**
* Worker-pool saturation: tasks in flight, threads, and jobs queued.
*/
inline constexpr char jobqSaturation[] = "jobq_saturation";
/**
* Jobs whose run time reached LoadMonitor's 1 s warn threshold, by job type.
* An exact counter for a process-wide freeze; the spans say what caused it.
*/
inline constexpr char jobqStallTotal[] = "jobq_stall_total";
// ===== Quorum and publish: can this node accept and publish a ledger? =======
/**
* Pre-accept quorum gate and publish lag, split by the `metric` label.
*/
inline constexpr char ledgerQuorumPublish[] = "ledger_quorum_publish";
/**
* Pre-accept gate rejections for being below quorum.
*/
inline constexpr char ledgerQuorumShortfallTotal[] = "ledger_quorum_shortfall_total";
// ===== Back-fill persistence: is history repair making progress? =============
/**
* Replay sub-acquires that fell back to a full ledger acquire, by stage.
*/
inline constexpr char ledgerReplayFallbackTotal[] = "ledger_replay_fallback_total";
/**
* Ledger replay tasks by terminal outcome.
*/
inline constexpr char ledgerReplayOutcomeTotal[] = "ledger_replay_outcome_total";
/**
* Forced jumps of the last closed ledger to a divergent chain.
*/
inline constexpr char ledgerJumpTotal[] = "ledger_jump_total";
// ===== Peer supply: what this node's peers can and will serve ===============
/**
* Peer coverage of the ledger sequence range this node still needs.
*/
inline constexpr char peerLedgerSupply[] = "peer_ledger_supply";
/**
* Inbound peer connection attempts, by terminal outcome.
*/
inline constexpr char peerAcceptTotal[] = "peer_accept_total";
/**
* Peer disconnects, by cause and connection direction.
*/
inline constexpr char peerDisconnectTotal[] = "peer_disconnect_total";
/**
* Peer data requests this node declined to serve, by kind and cause.
*/
inline constexpr char serveRefusedTotal[] = "serve_refused_total";
/**
* PeerFinder slots, dials in flight, and address-cache sizes.
*/
inline constexpr char peerfinderSlotCensus[] = "peerfinder_slot_census";
/**
* Amendment-block warning and the countdown to this node ceasing to validate.
*/
inline constexpr char amendmentBlock[] = "amendment_block";
// ===== Consensus =============================================================
/**
* Wall-clock duration of a completed consensus round.
*/
inline constexpr char consensusRoundDurationMs[] = "consensus_round_duration_ms";
/**
* Rounds where the network's preferred ledger differed from ours. One
* increment per transition into WrongLedger, labelled with the mode being
* left. WrongLedger itself is therefore never a value of the label.
*/
inline constexpr char consensusViewChangeTotal[] = "consensus_view_change_total";
// ===== Sweep: what the periodic cache sweep costs ============================
//
// The sweep runs every `SizedItem::SweepInterval` seconds (10 s on a tiny node
// through 120 s on a huge one), so these three emit sites are cold: one
// histogram Record and two counter Adds per sweep is free at that cadence.
/**
* Wall-clock duration of the `malloc_trim` call that ends every cache sweep.
*
* The cost of returning free heap pages to the kernel scales with the resident
* heap, so this is the signal that a node with a large existing database pays a
* per-sweep penalty a fresh node does not.
*/
inline constexpr char sweepMallocTrimUs[] = "sweep_malloc_trim_us";
/**
* Minor page faults taken *inside* the `malloc_trim` call.
*
* Cumulative, so `rate()` gives faults/sec. Scoped to the trim call only -- see
* the limitation noted on the runbook branch: this proves the trim itself
* faults, not that the trim causes later faults as the caches refill.
*/
inline constexpr char sweepMallocTrimMinorFaultsTotal[] = "sweep_malloc_trim_minor_faults_total";
/**
* Resident kilobytes the trim actually returned to the kernel.
*
* Cumulative and clamped at zero per sweep: a trim that reclaimed nothing, or
* during which another thread grew the heap faster than the trim shrank it,
* contributes 0 rather than a negative amount.
*/
inline constexpr char sweepMallocTrimReclaimedKbTotal[] = "sweep_malloc_trim_reclaimed_kb_total";
// ===== Rotation: the extra writes an online_delete rotation performs =========
/**
* Nodes re-stored by `copyNode` because they were missing from both backends.
*
* The genuinely unmeasured extra write of a rotation: a clean node reachable
* from the validated state map whose only on-disk copy lived in a backend an
* earlier rotation removed. Was warn-log-only.
*/
inline constexpr char rotationCopyNodeRestoreTotal[] = "rotation_copy_node_restore_total";
/**
* Rotation state: whether one is running, and the copy-forward write total.
*
* A gauge, not a counter, because the two readings are polled from the node
* store rather than pushed: `in_flight` is current state and `copy_forward` is a
* cumulative total the nodestore already keeps. Observed from the existing
* `registerNodeStoreGauge` callback, which is how everything else reads the node
* store from xrpld without libxrpl having to know about telemetry.
*/
inline constexpr char rotationState[] = "rotation_state";
/**
* Wall-clock seconds spent in one online-delete rotation phase. Labelled by
* `stage`; see lval::rotation_phase. Recorded once per phase end.
*/
inline constexpr char rotationPhaseDurationSeconds[] = "rotation_phase_duration_seconds";
// ===== Pre-existing instruments pulled in by the family ratchet ==============
//
// These predate the sync-diagnostics work. They are declared here because the
// checker's Rule I enforces literal-freedom per metric FAMILY (first
// underscore segment), and each of these shares a family with a name above --
// `ledger_`, `nodestore_`, `server_`, `peer_`, `state_`. Leaving them as
// literals would either weaken the rule to per-name (letting a typo'd sibling
// through) or require an exemption list. Declaring them is the honest option:
// no behaviour changes, and the next author editing these families finds the
// constant rather than inventing a second spelling.
//
// The remaining unconverted families are reported as Rule L warnings, so the
// outstanding work stays visible rather than silently accepted.
/**
* Built-vs-validated ledger mismatches, by reason.
*/
inline constexpr char ledgerHistoryMismatchTotal[] = "ledger_history_mismatch_total";
/**
* Ledger fee and economy readings.
*/
inline constexpr char ledgerEconomy[] = "ledger_economy";
/**
* NodeStore I/O counters, queue depth and write load.
*/
inline constexpr char nodestoreState[] = "nodestore_state";
/**
* Server-level health readings.
*/
inline constexpr char serverInfo[] = "server_info";
/**
* Peer-network quality readings.
*/
inline constexpr char peerQuality[] = "peer_quality";
/**
* Node state and operating-mode tracking.
*/
inline constexpr char stateTracking[] = "state_tracking";
} // namespace metric
/**
* Label keys -- the dimension names attached to a metric datapoint.
*
* Every key here is bounded by design: the values it can take are either a
* fixed set declared in `namespace lval` below, or a small enumeration the
* code derives (an operating mode, a job type). A key whose values are
* unbounded runtime data would mint one time series per distinct value, so no
* such key is declared.
*/
namespace label {
/**
* Sub-metric discriminator on a multi-series gauge.
*
* The pattern every observable gauge in MetricsRegistry.cpp already uses: one
* instrument carries several related readings, told apart by this label rather
* than by being separate instruments. Pre-dates the sync-diagnostics work;
* named here because the new gauges are its heaviest users.
*/
inline constexpr char metric[] = "metric";
/**
* Job type, as produced by `JobTypes::name()`.
*/
inline constexpr char jobType[] = "job_type";
/**
* Which producer submitted a job, within its job type.
*
* A job type has several producers (`RcvGetLedger` and `RcvGetObjByHash` both
* run as `JtLedgerReq`), so this is what attributes a latency spike to one of
* them. Bounded by `MetricsRegistry::sanitiseHandler()`, which folds any job
* name that is not all ASCII letters -- the ones embedding a ledger sequence
* -- down to a single `other` value.
*/
inline constexpr char handler[] = "handler";
/**
* Terminal result of a bounded operation.
*/
inline constexpr char outcome[] = "outcome";
/**
* Cause of a rejection, refusal or teardown.
*/
inline constexpr char reason[] = "reason";
/**
* Configured validator-list site URI. The one runtime-valued key here.
*/
inline constexpr char site[] = "site";
/**
* Operating mode a transition started from.
*/
inline constexpr char from[] = "from";
/**
* Operating mode a transition ended at.
*/
inline constexpr char to[] = "to";
/**
* Where acquired ledger data came from.
*/
inline constexpr char source[] = "source";
/**
* Which stage of a multi-step pipeline the event belongs to.
*/
inline constexpr char stage[] = "stage";
/**
* Connection direction, inbound or outbound.
*/
inline constexpr char direction[] = "direction";
/**
* Which kind of peer data request is being described.
*/
inline constexpr char request[] = "request";
/**
* Consensus mode being left when a view change is counted.
*/
inline constexpr char consensusMode[] = "consensus_mode";
} // namespace label
/**
* Bounded label values -- the fixed value sets the code itself writes.
*
* Nested by the instrument (or the gauge) that owns the set, because the same
* word means different things in different sets and a flat namespace would let
* two of them collide. A value is declared here only when the code chooses it
* from a fixed list; anything derived from runtime data stays out.
*/
namespace lval {
// ===== Shared outcome/direction slugs =======================================
/**
* Values shared by more than one instrument. Declared once so two instruments
* that mean the same thing cannot spell it differently.
*/
inline constexpr char timeout[] = "timeout";
inline constexpr char notFound[] = "not_found";
inline constexpr char inbound[] = "inbound";
inline constexpr char outbound[] = "outbound";
/**
* `dns_resolve_total` outcomes: did the resolver return any address?
*/
namespace dns_resolve {
inline constexpr char resolved[] = "resolved";
inline constexpr char empty[] = "empty";
} // namespace dns_resolve
/**
* `peer_accept_total` outcomes -- the nine exits of the inbound-accept path.
*
* Every exit records one of these, so the counter's total equals the number of
* inbound attempts and an unexplained gap is impossible.
*/
namespace peer_accept {
inline constexpr char localEndpointFail[] = "local_endpoint_fail";
inline constexpr char resourceLimit[] = "resource_limit";
inline constexpr char noSlot[] = "no_slot";
inline constexpr char notPeerRequest[] = "not_peer_request";
inline constexpr char protocolMismatch[] = "protocol_mismatch";
inline constexpr char badCookie[] = "bad_cookie";
inline constexpr char slotRefused[] = "slot_refused";
inline constexpr char accepted[] = "accepted";
inline constexpr char handshakeError[] = "handshake_error";
} // namespace peer_accept
/**
* `handshake_negotiation_fail_total` reasons -- one per rejection point in
* the handshake verifier.
*
* These separate a peer misconfiguration this node should tolerate
* (`wrong_network`, `self_connection`) from a local misconfiguration an
* operator must fix (`clock_skew`, `local_ip_mismatch`), which is the whole
* point of splitting the counter by reason.
*/
namespace handshake_fail {
inline constexpr char invalidServerDomain[] = "invalid_server_domain";
inline constexpr char invalidNetworkId[] = "invalid_network_id";
inline constexpr char wrongNetwork[] = "wrong_network";
inline constexpr char invalidClockTimestamp[] = "invalid_clock_timestamp";
inline constexpr char clockSkew[] = "clock_skew";
inline constexpr char unsupportedKeyType[] = "unsupported_key_type";
inline constexpr char badPublicKey[] = "bad_public_key";
inline constexpr char noSessionSignature[] = "no_session_signature";
inline constexpr char sessionVerifyFailed[] = "session_verify_failed";
inline constexpr char selfConnection[] = "self_connection";
inline constexpr char invalidLocalIp[] = "invalid_local_ip";
inline constexpr char localIpMismatch[] = "local_ip_mismatch";
inline constexpr char invalidRemoteIp[] = "invalid_remote_ip";
inline constexpr char remoteIpMismatch[] = "remote_ip_mismatch";
} // namespace handshake_fail
/**
* `unl_fetch_total` outcomes for the transport-level failures.
*
* The success path instead labels with `to_string(ListDisposition)`, whose
* values are owned by the protocol enum and are therefore not restated here --
* duplicating them would create a second place to update when a disposition is
* added.
*/
namespace unl_fetch {
inline constexpr char fetchError[] = "fetch_error";
inline constexpr char badStatus[] = "bad_status";
inline constexpr char parseError[] = "parse_error";
} // namespace unl_fetch
/**
* `unl_quorum` sub-metrics: the trusted-key count and the bar it must clear.
*/
namespace unl_quorum {
inline constexpr char trustedKeys[] = "trusted_keys";
inline constexpr char quorum[] = "quorum";
// 1 while the validator list has disabled quorum, 0 otherwise. A separate
// boolean rather than a sentinel value on `quorum` itself: a number that large
// plotted on a shared axis flattens the trusted-key line to the baseline,
// hiding the very failure it would mark.
inline constexpr char quorumDisabled[] = "quorum_disabled";
} // namespace unl_quorum
/**
* `clock_close_offset_seconds` sub-metric.
*/
namespace clock_offset {
inline constexpr char offset[] = "offset";
} // namespace clock_offset
/**
* `sync_state` sub-metrics -- the four "why am I not FULL" signals.
*
* `initial_full_duration_us` reads zero until the node first reaches FULL,
* which is the state this gauge exists to make visible rather than a missing
* value.
*/
namespace sync_state {
inline constexpr char initialFullDurationUs[] = "initial_full_duration_us";
inline constexpr char networkLedgerGate[] = "network_ledger_gate";
inline constexpr char serverStallSeconds[] = "server_stall_seconds";
inline constexpr char ledgersBehind[] = "ledgers_behind";
} // namespace sync_state
/**
* `sync_acquire` sub-metrics -- aggregate acquire progress.
*
* The two `missing_*_max` series are the stuck-detector: flat and non-zero
* across collection ticks means the acquire will never finish, while shrinking
* means slow but alive. `in_flight` is the context that tells idle from stuck.
*/
namespace sync_acquire {
inline constexpr char missingStateNodesMax[] = "missing_state_nodes_max";
inline constexpr char missingTxNodesMax[] = "missing_tx_nodes_max";
inline constexpr char receivedDataDepth[] = "received_data_depth";
inline constexpr char inFlight[] = "in_flight";
} // namespace sync_acquire
/**
* `shamap_cache_hit_rate` sub-metric: which cache the rate describes.
*/
namespace shamap_cache {
inline constexpr char treenode[] = "treenode";
} // namespace shamap_cache
/**
* `sync_acquire_source_total` sources: served locally or fetched from peers.
*/
namespace acquire_source {
inline constexpr char local[] = "local";
inline constexpr char network[] = "network";
} // namespace acquire_source
/**
* `sync_addnode_total` outcomes -- the per-node verdict on received SHAMap data.
*
* The split is what separates real progress (`good`) from wasted bandwidth
* (`duplicate`) and a misbehaving peer (`invalid`); traffic-level metrics show
* all three as healthy throughput.
*/
namespace addnode {
inline constexpr char good[] = "good";
inline constexpr char duplicate[] = "duplicate";
inline constexpr char invalid[] = "invalid";
} // namespace addnode
/**
* `jobq_saturation` sub-metrics: the numerator, denominator and the backlog.
*/
namespace jobq_saturation {
inline constexpr char runningTasks[] = "running_tasks";
inline constexpr char workerThreads[] = "worker_threads";
inline constexpr char totalWaiting[] = "total_waiting";
} // namespace jobq_saturation
/**
* `peer_ledger_supply` sub-metrics: who can serve what this node needs.
*/
namespace peer_supply {
inline constexpr char peersReporting[] = "peers_reporting";
inline constexpr char peersServingValidated[] = "peers_serving_validated";
inline constexpr char peersServingNext[] = "peers_serving_next";
inline constexpr char supplyMinSeq[] = "supply_min_seq";
inline constexpr char supplyMaxSeq[] = "supply_max_seq";
} // namespace peer_supply
/**
* `peerfinder_slot_census` sub-metrics -- slots, dials and address caches.
*
* `connecting` non-zero while `out_active` stays under `out_max` is the
* "starting dials and never completing them" case; both caches at zero on a
* fresh node means there is nothing left to dial at all.
*/
namespace slot_census {
inline constexpr char outActive[] = "out_active";
inline constexpr char outMax[] = "out_max";
inline constexpr char inActive[] = "in_active";
inline constexpr char inMax[] = "in_max";
inline constexpr char connecting[] = "connecting";
inline constexpr char fixedConfigured[] = "fixed_configured";
inline constexpr char fixedActive[] = "fixed_active";
inline constexpr char bootcache[] = "bootcache";
inline constexpr char livecache[] = "livecache";
} // namespace slot_census
/**
* `amendment_block` sub-metrics: the warning flag and the countdown.
*/
namespace amendment_block {
inline constexpr char warned[] = "warned";
inline constexpr char secondsToBlock[] = "seconds_to_block";
} // namespace amendment_block
/**
* `rotation_state` sub-metrics: is a rotation running, and how many extra
* writes have rotations caused.
*
* Read together: a `copy_forward` total that climbs while `in_flight` is 1 is
* the rotation doing its extra writes, which is the expected shape. The same
* total climbing while `in_flight` is 0 would mean the flag leaked, not that
* rotation is cheap.
*/
namespace rotation_state {
inline constexpr char inFlight[] = "in_flight";
inline constexpr char copyForward[] = "copy_forward";
} // namespace rotation_state
/**
* `stage` values for rotation_phase_duration_seconds. Identical to the
* child span suffixes in SHAMapStoreSpanNames.h so a panel can join the two.
*/
namespace rotation_phase {
inline constexpr char clearPrior[] = "clear_prior";
inline constexpr char copy[] = "copy";
inline constexpr char freshenKeys[] = "freshen.keys";
inline constexpr char freshenFetch[] = "freshen.fetch";
inline constexpr char newBackend[] = "new_backend";
inline constexpr char clearCaches[] = "clear_caches";
inline constexpr char swap[] = "swap";
inline constexpr char healthWait[] = "health_wait";
} // namespace rotation_phase
/**
* `metric` values added to the cache_metrics gauge for lock-hold peaks.
*/
namespace cache_metrics {
inline constexpr char treenodeLockHoldPeakUs[] = "treenode_lock_hold_peak_us";
inline constexpr char fullbelowLockHoldPeakUs[] = "fullbelow_lock_hold_peak_us";
} // namespace cache_metrics
/**
* `ledger_quorum_publish` sub-metrics: the gate, and how late publish is.
*/
namespace quorum_publish {
inline constexpr char trustedValidationTally[] = "trusted_validation_tally";
inline constexpr char quorumTarget[] = "quorum_target";
inline constexpr char timeToFirstValidatedUs[] = "time_to_first_validated_us";
inline constexpr char publishLag[] = "publish_lag";
} // namespace quorum_publish
/**
* `ledger_quorum_shortfall_total` stage: which gate did the rejecting.
*/
namespace quorum_shortfall {
inline constexpr char preAccept[] = "pre_accept";
} // namespace quorum_shortfall
/**
* `ledger_replay_fallback_total` stages: which sub-acquire gave up.
*/
namespace replay_fallback {
inline constexpr char skiplist[] = "skiplist";
inline constexpr char delta[] = "delta";
} // namespace replay_fallback
/**
* `ledger_replay_outcome_total` outcomes -- the four terminal states of a
* replay task. `timeout` is the shared slug above.
*/
namespace replay_outcome {
inline constexpr char success[] = "success";
inline constexpr char buildFailed[] = "build_failed";
inline constexpr char parameterFailed[] = "parameter_failed";
} // namespace replay_outcome
/**
* `peer_disconnect_total` reasons -- why a peer connection closed.
*
* The split separates our-fault backpressure (`large_sendq`,
* `charge_resources`) from a topology or network fault (`not_useful`,
* `ping_timeout`, `read_error`); the two call for opposite responses.
* `unknown` is the initial value and appears when a teardown path set no
* cause, so an unattributed disconnect is visible rather than absent.
*/
namespace disconnect {
inline constexpr char unknown[] = "unknown";
inline constexpr char malformedHandshake[] = "malformed_handshake";
inline constexpr char stopping[] = "stopping";
inline constexpr char chargeResources[] = "charge_resources";
inline constexpr char timerError[] = "timer_error";
inline constexpr char largeSendq[] = "large_sendq";
inline constexpr char notUseful[] = "not_useful";
inline constexpr char pingTimeout[] = "ping_timeout";
inline constexpr char shutdown[] = "shutdown";
inline constexpr char sharedValue[] = "shared_value";
inline constexpr char writeError[] = "write_error";
inline constexpr char graceful[] = "graceful";
inline constexpr char readError[] = "read_error";
} // namespace disconnect
/**
* `serve_refused_total` request kinds: what the peer had asked for.
*/
namespace serve_request {
inline constexpr char object[] = "object";
inline constexpr char fetchpack[] = "fetchpack";
inline constexpr char txset[] = "txset";
inline constexpr char ledger[] = "ledger";
} // namespace serve_request
/**
* `serve_refused_total` reasons: why this node would not answer.
* `not_found` is the shared slug above.
*
* `empty_reply` is the subtle one: the map WAS found, but the reply loop
* produced no nodes, so the requester still gets nothing and must ask another
* peer. Counting it as served would make a node that answers every request
* with an empty payload look healthy.
*/
namespace serve_refused {
inline constexpr char sendqFull[] = "sendq_full";
inline constexpr char loadShed[] = "load_shed";
inline constexpr char badType[] = "bad_type";
inline constexpr char noMap[] = "no_map";
inline constexpr char emptyReply[] = "empty_reply";
} // namespace serve_refused
} // namespace lval
} // namespace xrpl::telemetry

View File

@@ -0,0 +1,943 @@
#pragma once
/**
* Central OTel metrics registry: the export pipeline and the instruments that
* app code pushes values into.
*
* Owns the OpenTelemetry MeterProvider, the OTLP/HTTP exporter, the periodic
* reader and every SYNCHRONOUS instrument (counters and histograms) that is
* not already covered by the beast::insight StatsD pipeline. The instruments
* are created once at startup and drained by the OTel
* PeriodicExportingMetricReader at a fixed interval (10 s).
*
* When XRPL_ENABLE_TELEMETRY is **not** defined, this class compiles to a
* lightweight no-op: every public method is an empty inline.
*
* Every caller reaches it through ServiceRegistry::getMetricsRegistry(), and
* the XRPL_METRIC_* macros then create their own instruments from meter().
*
* Dependency / ownership diagram (ASCII):
*
* MetricsRegistry
* |
* +-- OTel MeterProvider (owns reader + exporter)
* | |
* | +-- PeriodicExportingMetricReader
* | +-- OtlpHttpMetricExporter
* |
* +-- Counters / Histograms (synchronous instruments)
* | +-- rpc_method_started_total
* | +-- rpc_method_finished_total
* | +-- rpc_method_errored_total
* | +-- rpc_method_us (Histogram)
* | +-- job_queued_total{job_type,handler}
* | +-- job_started_total{job_type,handler}
* | +-- job_finished_total{job_type,handler}
* | +-- job_queued_us{job_type,handler} (Histogram)
* | +-- job_running_us{job_type,handler} (Histogram)
* | +-- ledgers_closed_total
* | +-- validations_sent_total
* | +-- validations_checked_total
* | +-- ledger_history_mismatch_total{reason}
* | +-- txq_expired_total
* | +-- txq_dropped_total{reason}
* |
* +-- ValidationTracker (rolling validation-agreement windows)
*
* Control-flow for synchronous instruments:
*
* PerfLogImp::rpcStart/rpcEnd/jobQueue/jobStart/jobFinish
* |
* v
* MetricsRegistry::recordRpc*(method, ...) / recordJob*(type, ...)
* |
* v
* OTel Counter::Add() or Histogram::Record()
* |
* v
* Periodically flushed by the MetricReader
*
* Example usage:
*
* @code
* // In ApplicationImp's member-init list, right after telemetry_ and before
* // every subsystem. The constructor builds the pipeline and every
* // synchronous instrument, so no producer can exist before they do. The
* // endpoint, the TLS settings and the resource identity come from
* // [telemetry] and [network_id], read by Application.cpp rather than
* // through Telemetry::Setup.
* metricsRegistry_(std::make_unique<telemetry::MetricsRegistry>(
* telemetry_->isEnabled(), journal, options))
*
* // In PerfLogImp::rpcStart():
* if (auto* mr = app_.getMetricsRegistry())
* mr->recordRpcStarted("server_info");
*
* // In PerfLogImp::rpcEnd():
* if (auto* mr = app_.getMetricsRegistry())
* {
* mr->recordRpcFinished("server_info", durationUs);
* // or: mr->recordRpcErrored("server_info", durationUs);
* }
*
* // In PerfLogImp::jobQueue(). The second argument is the addJob name;
* // it is sanitised internally into the bounded `handler` label.
* if (auto* mr = app_.getMetricsRegistry())
* mr->recordJobQueued("ledgerData", "ProcessLData");
*
* // Shutdown, before any observer of live server state is torn down.
* // Idempotent, so run() and ~ApplicationImp both call it:
* metricsRegistry_->stop();
* @endcode
*
* Caveats:
* - The MetricsRegistry must be created AFTER the Telemetry object because
* it reads isEnabled() to decide whether to initialize the OTel SDK, and
* BEFORE every subsystem that records a metric. Declaration order in
* ApplicationImp is the guarantee; keep the member where it is.
* - Adding a new synchronous instrument requires updating both the header
* and the .cpp, then calling the new record*() method from the
* instrumentation site. Prefer the XRPL_METRIC_* macros, which need
* neither.
*/
#ifdef XRPL_ENABLE_TELEMETRY
// The tracker is held and exposed only in this configuration, where the
// observable-gauge callbacks that drain it exist.
#include <xrpl/telemetry/ValidationTracker.h>
#endif
#include <xrpl/beast/utility/Journal.h>
#include <algorithm>
#include <charconv>
#include <cstdint>
#include <limits>
#include <optional>
#include <string>
#include <string_view>
#include <system_error>
#include <utility>
#ifdef XRPL_ENABLE_TELEMETRY
#include <opentelemetry/metrics/meter.h>
#include <opentelemetry/metrics/meter_provider.h>
#include <opentelemetry/nostd/shared_ptr.h>
#include <opentelemetry/nostd/unique_ptr.h>
#include <opentelemetry/sdk/metrics/meter_provider.h>
// These two serve only the telemetry-only members below, so they are guarded
// like their uses: std::atomic by phase_, std::shared_ptr by provider_.
#include <atomic>
#include <memory>
#endif
namespace xrpl::telemetry {
/**
* Run time at which a finished job counts as a stall, in microseconds.
* Equal to LoadMonitor's 1 s warn threshold (LoadMonitor.cpp
* addLoadSample) so this counter and the "Job: ... run:" log line
* describe the same event.
*/
inline constexpr std::int64_t kJobStallThresholdUs = 1'000'000;
/**
* Central OpenTelemetry metric registry.
*
* Owns the metrics export pipeline and every push-model instrument that the
* beast::insight StatsD pipeline does not already cover. See the file-level
* header comment above for the instrument inventory and usage examples.
*
* Class / collaborator diagram (ASCII):
*
* MetricsRegistry
* |
* +-- creates/owns --> MeterProvider (SDK)
* | |
* | v
* | reader thread (~10 s) -> OTLP/HTTP export
* |
* +-- creates/owns --> Counter and Histogram instruments
* |
* +-- holds ----------> ValidationTracker (rolling windows)
*
* @note Thread safety:
* - The recordRpc, recordJob, and increment methods are invoked
* from hot paths. OTel Counter::Add() and Histogram::Record()
* are documented thread-safe, and null-guard checks protect
* uninitialized instruments.
* - recording() is a single acquire load and is read on every
* XRPL_METRIC_* call site, from any thread.
* - meter() may be called from any thread. The constructor is the
* last writer of the handle it returns; stop() leaves it alone.
* - ValidationTracker protects its rolling windows internally.
* - The constructor, hasPipeline() and stop() are NOT thread-safe
* with each other. All three read or write provider_, a plain
* shared_ptr that stop() resets, so all three belong on the
* single server lifecycle thread, in that order.
*
* @note Lifetime, in two phases (see Phase):
* - Ready: the constructor built the pipeline and the synchronous
* instruments. Runs in ApplicationImp's member-init list, so it precedes
* every subsystem that could record.
* - Stopped: stop() joined the reader thread. Runs before any observed
* service stops, from run() and again from ~ApplicationImp for the
* paths that never reach run().
*
* @note Extending:
* - Adding a new SYNCHRONOUS instrument (counter/histogram): prefer the
* XRPL_METRIC_* call-site macros in MetricMacros.h -- no header/cpp
* edit needed. Fall back to a dedicated member + init line + record
* method (the pattern below) only when the metric needs to be read
* back by other code (e.g. ValidationTracker-style accumulation) or
* needs a custom histogram bucket View (see the histogram note in
* MetricMacros.h).
* - An OBSERVABLE instrument does not belong here. Its callback reads live
* server state, so it must be registered only once that state exists,
* which is later than this object is built. Register it from the layer
* that owns those callbacks.
*/
class MetricsRegistry
{
public:
/**
* Everything the constructor needs from config: where to export, how to
* secure the connection, and the process identity stamped on the OTel
* resource.
*
* The values come from the `[telemetry]` section plus `[network_id]`, read
* by `makeMetricsRegistryOptions()` in `Application.cpp`. They must match
* what `makeTelemetrySetup()` gives the trace pipeline, or one node reports
* two identities and a dashboard filter shows half its series.
*
* A struct rather than ten positional parameters: seven of them are
* strings, so a swapped pair would compile and silently stamp the wrong
* label. Designated initializers name every value at the call site.
*
* @code
* MetricsRegistry::Options opts{
* .endpoint = "http://localhost:4318/v1/metrics",
* .serviceName = "xrpld",
* .serviceVersion = build_info::getVersionString(),
* .serviceInstanceId = nodePublicKey,
* .nodeId = nodePublicKey,
* .networkId = 2};
* MetricsRegistry registry(enabled, journal, opts);
*
* // Edge case: mutual TLS to a collector that requires it.
* opts.useTls = true;
* opts.tlsCaCertPath = "/etc/xrpld/otel-ca.pem";
* opts.tlsClientCertPath = "/etc/xrpld/node.pem";
* opts.tlsClientKeyPath = "/etc/xrpld/node.key";
* MetricsRegistry secure(enabled, journal, opts);
* @endcode
*
* @note Plain aggregate, no invariants enforced. `networkType` is not a
* field: it is derived from @ref networkId inside the constructor so
* the two can never disagree.
*/
struct Options
{
/**
* OTLP/HTTP endpoint URL for metric export, from
* `[telemetry] metrics_endpoint`.
*/
std::string endpoint;
/**
* service.name resource attribute, from `[telemetry] service_name`.
* Stamped unconditionally, so an empty value here yields an empty
* label rather than the SDK's `unknown_service` default. The caller
* seeds it with `systemName()`.
*/
std::string serviceName;
/**
* service.version resource attribute — the build's version string.
* Left off the resource when empty.
*/
std::string serviceVersion;
/**
* service.instance.id resource attribute, from
* `[telemetry] service_instance_id` or the node's base58 public key.
* Left off the resource when empty.
*/
std::string serviceInstanceId;
/**
* xrpl.node.id resource attribute — the node's base58 public key,
* which config cannot override. Left off the resource when empty.
*/
std::string nodeId;
/**
* Network identifier from `[network_id]`. Stamped as xrpl.network.id,
* and mapped to the xrpl.network.type label by `networkTypeFromId()`.
*/
std::uint32_t networkId{0};
/**
* Whether the exporter connects to the collector over TLS. The three
* paths below apply only when this is true.
*/
bool useTls{false};
/**
* CA bundle used to verify the collector. Empty selects the system
* CA store.
*/
std::string tlsCaCertPath;
/**
* This node's client certificate, presented for mutual TLS. Empty
* means one-way TLS.
*/
std::string tlsClientCertPath;
/**
* Private key for @ref tlsClientCertPath.
*/
std::string tlsClientKeyPath;
};
/**
* Construct the registry and, when enabled, build the whole metrics
* pipeline: OTLP exporter, periodic reader, MeterProvider and every
* SYNCHRONOUS instrument (counters and histograms).
*
* Doing this in the constructor is what fixes the init order. The
* Application declares its registry before every subsystem, so no
* producer can exist before the instruments do. A failure to build the
* pipeline is logged and leaves the registry a no-op; it never stops the
* node.
*
* @note Invariant for future changes: the constructor may create only
* instruments with NO callback of their own. Push-model counters
* and histograms qualify; app code records into them when it is
* ready. An instrument registered here is live immediately, and
* the reader thread may invoke its callback before the rest of
* the server is built, so any observable whose callback reads
* live server state must be registered later, by the layer that
* owns those callbacks. This applies to observable COUNTERS as
* well as gauges.
*
* @param enabled False makes every method a no-op (telemetry disabled).
* @param journal Log output.
* @param options Endpoint, TLS settings and resource identity, all read
* from config by the caller. See @ref Options.
*/
MetricsRegistry(bool enabled, beast::Journal journal, Options const& options);
/**
* Stops the pipeline if run() or ~ApplicationImp did not already.
*/
~MetricsRegistry();
/**
* Non-copyable, non-movable.
*/
MetricsRegistry(MetricsRegistry const&) = delete;
MetricsRegistry&
operator=(MetricsRegistry const&) = delete;
/**
* Flush pending metrics and shut down the pipeline.
*
* Stores `Phase::Stopped` first so `recording()` reads false on every
* later record call, then destroys the SDK provider. meter_ is not
* touched: record threads may still be running, and the gate is what
* keeps them off the dying pipeline. Idempotent.
*
* @pre Anything that observes live server state on the reader thread has
* already been disarmed. Shutting the provider down joins that
* thread, so a caller that has not disarmed its observers leaves a
* narrow race between the final tick and the teardown of what those
* observers read.
*/
void
stop();
/**
* @return true if the registry is actively exporting metrics.
*/
[[nodiscard]] bool
isEnabled() const noexcept
{
return enabled_;
}
/**
* @return true when a record call is safe to run.
*
* False when the registry is disabled, or after stop() has torn down the
* export pipeline. After stop() the SDK's SyncMetricStorage still holds a
* raw pointer to an AggregationConfig owned by a destroyed View, so a
* record with a first-seen attribute set would fire the factory lambda
* and deref that dangling pointer. Every XRPL_METRIC_* macro reads this
* once before touching an instrument.
*
* One acquire atomic load in the hot path.
*/
[[nodiscard]] bool
recording() const noexcept
{
#ifdef XRPL_ENABLE_TELEMETRY
return enabled_ && phase_.load(std::memory_order_acquire) != Phase::Stopped;
#else
return enabled_;
#endif
}
/**
* @return true when a real exporting pipeline exists, as opposed to the
* no-op meter installed when the pipeline is disabled.
*
* A meter() check cannot answer this. The registry always hands out a
* meter, so registering instruments on a no-op one would report success
* and export nothing. Ask this before registering an observable
* instrument.
*
* @note Not thread-safe against stop(), which drops the provider this
* reads. Call it from the server lifecycle thread, like the constructor
* and stop().
*/
[[nodiscard]] bool
hasPipeline() const noexcept;
// -----------------------------------------------------------------
// Synchronous instrument recording (called from PerfLog hot paths)
// -----------------------------------------------------------------
/**
* Record an RPC method call start.
* @param method The RPC method name (e.g. "server_info").
*/
void
recordRpcStarted(std::string_view method);
/**
* Record an RPC method call completion.
* @param method The RPC method name.
* @param durationUs Execution time in microseconds.
*/
void
recordRpcFinished(std::string_view method, std::int64_t durationUs);
/**
* Record an RPC method call error.
* @param method The RPC method name.
* @param durationUs Execution time in microseconds.
*/
void
recordRpcErrored(std::string_view method, std::int64_t durationUs);
/**
* The `handler` label value used for any job name that fails the
* sanitiser's all-ASCII-letters rule.
*
* Public because both sanitiseHandler() and its unit tests must agree
* on the exact fallback token; a test asserting against its own copy
* of the string would not catch a change made here.
*
* Declared as std::string_view rather than the `constexpr char k[]`
* form used for instrument names in MetricsRegistry.cpp: this value is
* *returned* by sanitiseHandler(), whose return type is
* std::string_view, and is compared against std::string_view in tests.
* Matching the type avoids array-to-pointer decay and a needless
* strlen at each use.
*/
static constexpr std::string_view kHandlerOther{"other"};
/**
* Reduce a job name to a bounded-cardinality `handler` label value.
*
* A job type can have several producers — both `RcvGetLedger` and
* `RcvGetObjByHash` run as `JtLedgerReq` — so `job_type` alone cannot
* attribute a latency spike to one of them. The job name can, but it
* cannot be used raw: two names embed a ledger sequence number
* (`"Pub" + std::to_string(seq)` in LedgerPersistence.cpp and
* `"OB" + std::to_string(...)` in OrderBookDBImpl.cpp), which would
* mint a fresh Prometheus series for every ledger.
*
* The rule is therefore: keep the name only when it is non-empty and
* every character is an ASCII letter; otherwise return `"other"`.
* Both dynamic names always contain digits, so they always fold to
* `"other"`, while every all-letter name is a compile-time literal.
* The label domain is thus a function of the literals present in the
* source — 43 names plus `"other"` at the time of writing — and
* cannot grow at runtime. A name added later that does not satisfy
* the rule degrades to `"other"` rather than becoming unbounded,
* which is a stronger guarantee than an allowlist that would have to
* be maintained by hand.
*
* Defined inline so it is available in a build without telemetry and
* usable in a constant expression.
*
* @param name The job name as passed to JobQueue::addJob.
* @return @p name when it is non-empty and all ASCII letters, else
* kHandlerOther.
*
* @note Pure and reentrant: holds no state, performs no I/O, and is
* safe to call concurrently from any thread.
* @note The letter test is an explicit ASCII range check rather than
* std::isalpha, which classifies by the current C locale. A
* locale-dependent test could admit non-ASCII bytes and so
* weaken the cardinality bound this function exists to provide.
* @note When the name is kept, the returned view aliases @p name, so
* it must not outlive the caller's buffer. The kHandlerOther case
* returns a view of a static constant and is always valid.
*
* Example:
* @code
* sanitiseHandler("RcvGetObjByHash"); // "RcvGetObjByHash"
* sanitiseHandler("Pub94512331"); // kHandlerOther (digits)
* sanitiseHandler(""); // kHandlerOther (empty)
* @endcode
*/
[[nodiscard]] static constexpr std::string_view
sanitiseHandler(std::string_view name) noexcept
{
auto const isAsciiLetter = [](char const c) {
return (c >= 'a' && c <= 'z') || (c >= 'A' && c <= 'Z');
};
if (name.empty() || !std::ranges::all_of(name, isAsciiLetter))
return kHandlerOther;
return name;
}
/**
* Divide a cumulative total by its count, optionally scaled, reporting
* absence rather than zero when the count is zero.
*
* Every cumulative counter this registry publishes has a companion mean
* that is only defined once the counter has moved. Reporting such a mean
* as `0` is worse than not reporting it: `0` is a plausible reading, so a
* dashboard draws a flat line at the bottom of the axis and an operator
* concludes "reads are instant" when the truth is "nothing has been
* read". Returning std::nullopt makes the caller skip the observation, so
* the series has a genuine gap instead.
*
* @p scale exists because the gauge these feed is integral. A mean writer
* depth of 1.4 truncates to 1, which is indistinguishable from a healthy
* 1.0, so the caller scales by 100 and says so in the metric name.
*
* The arithmetic divides before scaling and scales the remainder
* separately, so a long-lived node cannot overflow the product. Should
* the result still exceed the gauge's range it saturates at
* INT64_MAX rather than wrapping, because a wrapped gauge reads as a
* sudden healthy-looking dip.
*
* Defined inline for the same reason as sanitiseHandler(); constexpr so
* the cases below are checked at compile time.
*
* @param total Cumulative numerator (e.g. summed microseconds).
* @param count Number of samples in @p total.
* @param scale Fixed-point multiplier applied to the quotient. Must be
* at least 1; 0 is meaningless and yields std::nullopt.
* @return The scaled mean, or std::nullopt when @p count is 0 (mean
* undefined) or @p scale is 0.
*
* @note Pure and reentrant: holds no state and performs no I/O.
* @note Truncates toward zero, like integer division. A mean of 9.9 us
* reads as 9 at @p scale 1 and as 990 at @p scale 100.
*
* Example:
* @code
* scaledMean(500, 4); // 125 -- mean microseconds
* scaledMean(7, 5, 100); // 140 -- mean 1.4, scaled by 100
* scaledMean(500, 0); // nullopt -- no samples, so no mean
* @endcode
*/
[[nodiscard]] static constexpr std::optional<std::int64_t>
scaledMean(std::uint64_t total, std::uint64_t count, std::uint64_t scale = 1) noexcept
{
if (count == 0 || scale == 0)
return std::nullopt;
constexpr auto kInt64Max =
static_cast<std::uint64_t>(std::numeric_limits<std::int64_t>::max());
auto const whole = total / count;
if (whole > kInt64Max / scale)
return static_cast<std::int64_t>(kInt64Max);
// Scale the remainder too, so `scale` recovers the fractional digits
// it exists for. Skipped when the product itself would overflow, at
// which point it is worth less than one part in 2^63 of the result.
auto const remainder = total % count;
std::uint64_t fraction = 0;
if (remainder <= std::numeric_limits<std::uint64_t>::max() / scale)
fraction = remainder * scale / count;
auto const scaled = whole * scale;
if (scaled > kInt64Max - fraction)
return static_cast<std::int64_t>(kInt64Max);
return static_cast<std::int64_t>(scaled + fraction);
}
/**
* Read one comma-separated segment of a complete-ledger range string.
*
* The producer is xrpl::to_string(RangeSet), documented in
* xrpl/basics/RangeSet.h. It renders an interval as `first-last`, and an
* interval whose first equals its last as a bare sequence number. A segment
* with no dash is therefore a range of one ledger, not a malformed one.
*
* Defined inline for the same reason as sanitiseHandler().
*
* @param segment One segment, already split on ','. Leading or trailing
* whitespace is rejected, because the producer emits none.
* @return The inclusive first and last sequence of the range. The two are
* equal for a single-ledger range. std::nullopt when @p segment is not
* something this producer can emit.
*
* @note Pure and reentrant: holds no state, performs no I/O, and is safe to
* call concurrently from any thread.
* @note Reports malformed input instead of throwing, so one unreadable
* segment costs its own range and not every range after it.
* @note A reversed range such as "9-4" is returned as given. RangeSet
* cannot emit one.
*
* Example:
* @code
* parseLedgerRange("32570-50000"); // {32570, 50000}
* parseLedgerRange("5000"); // {5000, 5000} -- one ledger
* parseLedgerRange("5-"); // nullopt
* @endcode
*/
[[nodiscard]] static std::optional<std::pair<std::uint32_t, std::uint32_t>>
parseLedgerRange(std::string_view segment) noexcept
{
auto const parseSeq = [](std::string_view text) -> std::optional<std::uint32_t> {
std::uint32_t value = 0;
auto const* const begin = text.data();
auto const* const end = begin + text.size();
auto const [ptr, ec] = std::from_chars(begin, end, value);
// from_chars stops at the first character it cannot use, so the
// whole segment counts as read only when it consumed all of it.
if (ec != std::errc{} || ptr != end)
return std::nullopt;
return value;
};
auto const dash = segment.find('-');
if (dash == std::string_view::npos)
{
auto const only = parseSeq(segment);
if (!only)
return std::nullopt;
return std::pair{*only, *only};
}
auto const first = parseSeq(segment.substr(0, dash));
auto const last = parseSeq(segment.substr(dash + 1));
if (!first || !last)
return std::nullopt;
return std::pair{*first, *last};
}
/**
* Record a job enqueued event.
* @param jobType The job type name (e.g. "ledgerData").
* @param jobName The addJob name, reduced to a bounded `handler`
* label by sanitiseHandler(). Distinguishes producers
* that share a job type.
*/
void
recordJobQueued(std::string_view jobType, std::string_view jobName);
/**
* Record a job start event.
* @param jobType The job type name.
* @param jobName The addJob name; see recordJobQueued().
* @param queuedDurUs Time the job spent waiting in the queue (us).
*/
void
recordJobStarted(std::string_view jobType, std::string_view jobName, std::int64_t queuedDurUs);
/**
* Record a job finish event.
* @param jobType The job type name.
* @param jobName The addJob name; see recordJobQueued().
* @param runningDurUs Execution time in microseconds.
*/
void
recordJobFinished(
std::string_view jobType,
std::string_view jobName,
std::int64_t runningDurUs);
// -----------------------------------------------------------------
// External dashboard parity counters
// -----------------------------------------------------------------
/**
* Increment the ledgers_closed_total counter.
*
* @note Currently has no callers: the ledgers_closed_total counter is
* incremented at its consensus call site via the XRPL_METRIC_COUNTER_INC
* macro (see MetricMacros.h). This method and its eagerly-created
* counter are retained as a fallback and are slated for removal in a
* separate cleanup once the macro path has proven out.
*/
void
incrementLedgersClosed();
/**
* Increment the validations_sent_total counter.
* Called from RCLConsensus::Adaptor::validate() when a validation
* is produced and broadcast.
*/
void
incrementValidationsSent();
/**
* Increment the validations_checked_total counter.
* Called from NetworkOPs::recvValidation() when a network validation
* is received and checked.
*/
void
incrementValidationsChecked();
/**
* Increment the ledger_history_mismatch_total counter for a reason.
* Called from LedgerHistory::handleMismatch() once the mismatch has
* been classified. The reason label turns fork diagnosis from a
* log-grep into a queryable time series.
* @param reason Classified mismatch cause (e.g. "prior_ledger",
* "close_time", "consensus_txset", "same_txset_diff_result",
* "unknown").
*/
void
incrementLedgerHistoryMismatch(std::string_view reason);
/**
* Increment the txq_expired_total counter.
* Called from TxQ::processClosedLedger() for each queued transaction
* removed because its LastLedgerSequence has passed — submitters who
* under-bid the escalating fee and were never included.
*/
void
incrementTxqExpired();
/**
* Increment the txq_dropped_total{reason} counter.
* Called from TxQ::apply() when a transaction is refused admission to
* the queue (e.g. the queue is full). Distinct from expiry (already
* queued) and from jq_trans_overflow (job queue, not TxQ).
* @param reason Admission-control rejection cause (e.g. "queue_full").
*/
void
incrementTxqDropped(std::string_view reason);
#ifdef XRPL_ENABLE_TELEMETRY
/**
* Access the validation agreement tracker.
* Used by consensus and ledger hooks to record our validations and
* network validations so the tracker can compute agreement percentages.
*
* Guarded, along with the tracker itself, because only the observable-gauge
* callbacks read it and those exist only in this configuration. Recording
* into it is not free: each call takes its lock and inserts an entry.
* @return Reference to the internal ValidationTracker instance.
*/
[[nodiscard]] ValidationTracker&
getValidationTracker()
{
return validationTracker_;
}
/**
* Access the shared OTel Meter for call-site instrument creation.
* Used by the XRPL_METRIC_* macros (MetricMacros.h) so new synchronous
* counters/histograms can be declared at their call site instead of as
* MetricsRegistry members.
*
* Invariant: never empty while recording() is true. The constructor sets
* it to the real meter, or to a no-op meter when the pipeline failed to
* build, and never writes it again, so reads need no lock. After stop()
* the meter's SDK context is gone; the macros gate on recording() first,
* so no caller reaches it then.
*
* @return The shared Meter.
*/
[[nodiscard]] opentelemetry::nostd::shared_ptr<opentelemetry::metrics::Meter>
meter() const noexcept
{
return meter_;
}
#endif
private:
/**
* Master enable flag; when false all methods are no-ops.
*/
bool const enabled_;
#ifdef XRPL_ENABLE_TELEMETRY
/**
* Tracks validation agreement between this node and the network.
*
* Guarded because reconcile() -- which resolves and then prunes recorded
* events -- runs only from the observable-gauge callbacks. Recording
* without it accumulates one entry per validated ledger, so the tracker
* exists only where something drains it.
*/
ValidationTracker validationTracker_;
/**
* Journal for logging.
*/
beast::Journal const journal_;
/**
* Where the registry is in its life. Construction ends in `Ready`;
* stop() moves to `Stopped`.
*
* After `Stopped` the SDK pipeline is gone. recording() reads false, so
* no macro touches meter_ or a cached instrument.
*/
enum class Phase { Ready, Stopped };
/**
* Current phase. Written from the server lifecycle thread with release
* ordering; read from record threads via `recording()` with acquire
* ordering, so no record starts once stop() has stored `Stopped`.
*/
std::atomic<Phase> phase_{Phase::Ready};
/**
* The SDK MeterProvider that owns the export pipeline.
*/
std::shared_ptr<opentelemetry::sdk::metrics::MeterProvider> provider_;
/**
* The Meter used to create all instruments.
*/
opentelemetry::nostd::shared_ptr<opentelemetry::metrics::Meter> meter_;
// --- Synchronous instruments (RPC) ---
/**
* Counter: rpc_method_started_total{method="<name>"}
*/
opentelemetry::nostd::unique_ptr<opentelemetry::metrics::Counter<uint64_t>> rpcStartedCounter_;
/**
* Counter: rpc_method_finished_total{method="<name>"}
*/
opentelemetry::nostd::unique_ptr<opentelemetry::metrics::Counter<uint64_t>> rpcFinishedCounter_;
/**
* Counter: rpc_method_errored_total{method="<name>"}
*/
opentelemetry::nostd::unique_ptr<opentelemetry::metrics::Counter<uint64_t>> rpcErroredCounter_;
/**
* Histogram: rpc_method_us{method="<name>"}
*/
opentelemetry::nostd::unique_ptr<opentelemetry::metrics::Histogram<double>>
rpcDurationHistogram_;
// --- Synchronous instruments (Job Queue) ---
// All five carry handler="<sanitised addJob name>" in addition to
// job_type, so producers that share a job type stay distinguishable.
/**
* Counter: job_queued_total{job_type="<name>",handler="<name>"}
*/
opentelemetry::nostd::unique_ptr<opentelemetry::metrics::Counter<uint64_t>> jobQueuedCounter_;
/**
* Counter: job_started_total{job_type="<name>",handler="<name>"}
*/
opentelemetry::nostd::unique_ptr<opentelemetry::metrics::Counter<uint64_t>> jobStartedCounter_;
/**
* Counter: job_finished_total{job_type="<name>",handler="<name>"}
*/
opentelemetry::nostd::unique_ptr<opentelemetry::metrics::Counter<uint64_t>> jobFinishedCounter_;
/**
* Counter: jobq_stall_total{job_type="<name>"} — one per finished job
* whose run time reached kJobStallThresholdUs.
*/
opentelemetry::nostd::unique_ptr<opentelemetry::metrics::Counter<uint64_t>> jobStallCounter_;
/**
* Histogram: job_queued_us{job_type="<name>",handler="<name>"}
*/
opentelemetry::nostd::unique_ptr<opentelemetry::metrics::Histogram<double>>
jobQueuedDurationHistogram_;
/**
* Histogram: job_running_us{job_type="<name>",handler="<name>"}
*/
opentelemetry::nostd::unique_ptr<opentelemetry::metrics::Histogram<double>>
jobRunningDurationHistogram_;
// --- External dashboard parity counters ---
/**
* Counter: ledgers_closed_total — incremented each consensus round.
*/
opentelemetry::nostd::unique_ptr<opentelemetry::metrics::Counter<uint64_t>>
ledgersClosedCounter_;
/**
* Counter: validations_sent_total — incremented when this node sends a validation.
*/
opentelemetry::nostd::unique_ptr<opentelemetry::metrics::Counter<uint64_t>>
validationsSentCounter_;
/**
* Counter: validations_checked_total — incremented for each network validation
* received.
*/
opentelemetry::nostd::unique_ptr<opentelemetry::metrics::Counter<uint64_t>>
validationsCheckedCounter_;
/**
* Counter: ledger_history_mismatch_total{reason} — incremented per classified
* built-vs-validated ledger mismatch.
*/
opentelemetry::nostd::unique_ptr<opentelemetry::metrics::Counter<uint64_t>>
ledgerHistoryMismatchCounter_;
/**
* Counter: txq_expired_total — incremented per transaction expired out of the
* transaction queue.
*/
opentelemetry::nostd::unique_ptr<opentelemetry::metrics::Counter<uint64_t>> txqExpiredCounter_;
/**
* Counter: txq_dropped_total{reason} — incremented when a transaction is refused
* admission to the queue.
*/
opentelemetry::nostd::unique_ptr<opentelemetry::metrics::Counter<uint64_t>> txqDroppedCounter_;
/**
* Build the OTLP/HTTP exporter, periodic reader, resource attributes and
* histogram views, then create the MeterProvider and meter. Extracted
* from the constructor to keep each function under the 80-line limit.
*
* @param options Endpoint, TLS settings and resource identity, forwarded
* unchanged from the constructor. See @ref Options.
*/
void
initExporterAndProvider(Options const& options);
/**
* Create the synchronous instruments (RPC and job-queue counters and
* histograms, plus the external dashboard parity counters). Extracted
* from the constructor to keep each function under the 80-line limit.
*/
void
initSyncInstruments();
/**
* Give up the pipeline after a build failure: drop the provider, hand
* out a no-op meter so every call site still gets an instrument, and log
* why. The registry stays enabled and inert for the process.
*
* @param reason What failed, for the log line.
*/
void
disablePipeline(std::string_view reason);
#endif // XRPL_ENABLE_TELEMETRY
};
} // namespace xrpl::telemetry

View File

@@ -0,0 +1,758 @@
#pragma once
/**
* @file ValidationTracker.h
* Standalone validation agreement tracker for telemetry.
*/
#include <xrpl/basics/UnorderedContainers.h>
#include <xrpl/basics/base_uint.h>
#include <xrpl/protocol/Protocol.h>
#include <boost/smart_ptr/atomic_shared_ptr.hpp>
#include <boost/smart_ptr/shared_ptr.hpp>
#include <array>
#include <atomic>
#include <chrono>
#include <cstddef>
#include <cstdint>
#include <deque>
namespace xrpl::telemetry {
/**
* Tracks whether this validator's validations agree with network consensus,
* over rolling 1-hour, 24-hour and 7-day windows plus lifetime totals.
*
* Two independent events are recorded per ledger:
* 1. "We validated" -- our node published a validation for a ledger hash.
* 2. "Network validated" -- the network reached consensus on that hash.
*
* reconcile() compares the two flags once the grace period has passed. Both
* flags set is an agreement, anything else is a miss. A miss becomes an
* agreement if its other half arrives inside the late-repair window.
*
* The writer paths take no lock at all: each writer owns one ring, so the two
* never touch shared state. Everything a decision needs belongs to whichever
* thread is inside reconcile(), and a second caller returns instead of waiting.
* Readers take a copy of the snapshot the reducer published.
*
* Data flow:
* @code
* RCLConsensus::Adaptor::validate LedgerMaster::setValidLedger
* | |
* recordOurValidation() recordNetworkValidation()
* v v
* +------------+ +--------------+
* | ourRing_ | | networkRing_ |
* +------------+ +--------------+
* \ /
* \--------- reconcile() ------------/
* |
* pending_ --> one-minute buckets --> 1h / 24h / 7d counters
* |
* published_ snapshot
* |
* agreementPct1h() / agreements24h() / missed7d() / ...
* @endcode
*
* Usage -- basic recording and querying:
* @code
* xrpl::telemetry::ValidationTracker tracker;
*
* // On local validation:
* tracker.recordOurValidation(ledgerHash, seq);
*
* // On network consensus:
* tracker.recordNetworkValidation(ledgerHash, seq);
*
* // Periodically (e.g. every ten seconds):
* tracker.reconcile();
*
* // Query agreement percentage:
* double pct = tracker.agreementPct1h();
* @endcode
*
* Usage -- edge case with late arrival:
* @code
* xrpl::telemetry::ValidationTracker tracker;
*
* // Network validates first, our validation arrives late:
* tracker.recordNetworkValidation(hash, seq);
* tracker.reconcile(); // counted as a miss
*
* // Our validation arrives inside the repair window:
* tracker.recordOurValidation(hash, seq);
* tracker.reconcile(); // repaired to an agreement
* @endcode
*
* Usage -- a test drives the clock so a window edge is reachable at once:
* @code
* // A lambda with no captures converts to the function pointer NowFn wants.
* static xrpl::telemetry::ValidationTracker::TimePoint fakeNow{};
* xrpl::telemetry::ValidationTracker tracker([] { return fakeNow; });
*
* tracker.recordOurValidation(hash, seq);
* fakeNow += xrpl::telemetry::ValidationTracker::gracePeriod();
* tracker.reconcile(); // decides the event without any waiting
* @endcode
*
* @note Thread-safety: every public method may be called concurrently. The
* two record methods each need a single writer thread, which is how consensus
* and the ledger master call them. reconcile() and the getters may be called
* from any thread and any number of threads.
* @note reconcile() and the getters share the published snapshot through an
* atomic shared_ptr, which every implementation guards with a short internal
* spin. No writer path touches it, so nothing a consensus thread calls can
* spin. Boost's is used because Apple's libc++ has no std::atomic for a
* shared_ptr, so the std spelling does not compile there.
* @note A writer whose ring is full discards the event and bumps
* droppedEvents(). Counts are then low but never wrong.
* @note Window edges are rounded to whole minutes, because counts are kept in
* one-minute buckets.
*/
class ValidationTracker
{
public:
/**
* Monotonic clock used for all internal timestamps.
*/
using Clock = std::chrono::steady_clock;
/**
* Time point type from the monotonic clock.
*/
using TimePoint = Clock::time_point;
/**
* Time source. A test supplies its own so it can reach a window edge
* without waiting for one.
*/
using NowFn = TimePoint (*)();
/**
* Construct a tracker reading time from the given source.
* @param now Function returning the current time. Defaults to the
* monotonic clock, so default construction works.
*/
explicit ValidationTracker(NowFn now = &ValidationTracker::steadyNow) : now_(now)
{
}
/**
* Record that this node sent a validation for the given ledger.
* @param ledgerHash Hash of the ledger we validated.
* @param seq Ledger sequence number.
*/
void
recordOurValidation(uint256 const& ledgerHash, LedgerIndex seq);
/**
* Record that the network reached consensus on the given ledger.
* @param ledgerHash Hash of the network-validated ledger.
* @param seq Ledger sequence number.
*/
void
recordNetworkValidation(uint256 const& ledgerHash, LedgerIndex seq);
/**
* Drain both rings, decide every event past the grace period, retire
* expired buckets and publish a fresh snapshot, in that order.
*
* Call periodically, for example every ten seconds. Returns without doing
* anything if another thread is already inside, so a losing caller reads
* data at most one cycle old rather than blocking.
*/
void
reconcile();
/**
* @name Rolling-window percentage getters
*/
/** @{ */
/**
* Agreement percentage over the last 1 hour.
* @return Percentage [0.0, 100.0], or 0.0 if no data.
*/
[[nodiscard]] double
agreementPct1h() const;
/**
* Agreement percentage over the last 24 hours.
* @return Percentage [0.0, 100.0], or 0.0 if no data.
*/
[[nodiscard]] double
agreementPct24h() const;
/**
* Agreement percentage over the last 7 days.
* @return Percentage [0.0, 100.0], or 0.0 if no data.
*/
[[nodiscard]] double
agreementPct7d() const;
/** @} */
/**
* @name Rolling-window count getters
*/
/** @{ */
/**
* Number of agreements in the 1-hour window.
*/
[[nodiscard]] std::uint64_t
agreements1h() const;
/**
* Number of misses in the 1-hour window.
*/
[[nodiscard]] std::uint64_t
missed1h() const;
/**
* Number of agreements in the 24-hour window.
*/
[[nodiscard]] std::uint64_t
agreements24h() const;
/**
* Number of misses in the 24-hour window.
*/
[[nodiscard]] std::uint64_t
missed24h() const;
/**
* Number of agreements in the 7-day window.
*/
[[nodiscard]] std::uint64_t
agreements7d() const;
/**
* Number of misses in the 7-day window.
*/
[[nodiscard]] std::uint64_t
missed7d() const;
/** @} */
/**
* @name Lifetime totals (atomic, lock-free reads)
*/
/** @{ */
/**
* Total agreements since process start.
*/
[[nodiscard]] std::uint64_t
totalAgreements() const;
/**
* Total misses since process start.
*/
[[nodiscard]] std::uint64_t
totalMissed() const;
/**
* Agreements counted at first classification, never adjusted afterwards.
*
* @note Unlike totalAgreements(), this only ever rises. A late repair
* leaves it alone, so it can back a Prometheus counter, which must never
* decrease.
*/
[[nodiscard]] std::uint64_t
totalAgreementsEver() const;
/**
* Misses counted at first classification, never adjusted afterwards.
*
* @note Unlike totalMissed(), this only ever rises. A miss that a late
* repair turns into an agreement stays counted here.
*/
[[nodiscard]] std::uint64_t
totalMissedEver() const;
/**
* Total validations this node sent.
*/
[[nodiscard]] std::uint64_t
totalValidationsSent() const;
/**
* Total network validations observed for comparison.
*/
[[nodiscard]] std::uint64_t
totalValidationsChecked() const;
/**
* Events a writer discarded because its ring was full.
* @return Lifetime count of discards across both rings. Non-zero means
* reconcile() is not being called often enough.
*/
[[nodiscard]] std::uint64_t
droppedEvents() const;
/** @} */
/**
* @name Bounds, exposed so a test states the same bound as the code
*/
/** @{ */
/**
* Slots in each writer's ring.
*/
static constexpr std::size_t
ringCapacity()
{
return kRingCapacity;
}
/**
* Delay before an event is decided, so both sides can arrive first.
*/
static constexpr std::chrono::seconds
gracePeriod()
{
return kGracePeriod;
}
/**
* How long after a decision a miss can still become an agreement.
*/
static constexpr std::chrono::minutes
lateRepairWindow()
{
return kLateRepairWindow;
}
/** @} */
private:
/**
* Slots per ring. A power of two, so the index is a mask rather than a
* division. 128 slots is about 8.5 minutes of ledgers at one every four
* seconds, against a normal drain gap of well under a minute.
*/
static constexpr std::size_t kRingCapacity = 128;
/**
* Grace period before deciding a ledger event.
*/
static constexpr auto kGracePeriod = std::chrono::seconds(8);
/**
* Window during which a missed event can be repaired.
*/
static constexpr auto kLateRepairWindow = std::chrono::minutes(5);
/**
* One-minute buckets spanned by the short window.
*/
static constexpr std::size_t kBuckets1h = 60;
/**
* One-minute buckets spanned by the long window.
*/
static constexpr std::size_t kBuckets24h = 24 * 60;
/**
* One-minute buckets spanned by the extended window. Also the length of
* the bucket grid, so a slot is reused exactly seven days later.
*/
static constexpr std::size_t kBuckets7d = 7 * 24 * 60;
/**
* Maximum number of ledger hashes remembered as already counted.
* At one ledger every four seconds this spans about eleven hours.
* A validation arriving for a ledger counted before that is counted
* again.
*/
static constexpr std::size_t kMaxTalliedEvents = 10000;
/**
* Default time source.
* @return The monotonic clock's current time point.
*/
static TimePoint
steadyNow();
/**
* One recorded event as it travels from a writer to the reducer.
*/
struct Slot
{
uint256 hash; ///< Ledger hash being reported.
/**
* Ledger sequence number as the caller gave it. Carried for
* diagnostics: the counters key on the hash, not the sequence.
*/
LedgerIndex seq{0};
TimePoint at; ///< When the writer recorded it.
};
/**
* Single-producer, single-consumer ring of recorded events.
*
* The producer only advances head_ and the consumer only advances tail_,
* so a release store on one side and an acquire load on the other is
* enough: no compare-exchange, no retry loop, nothing to block on.
*
* @code
* Ring r;
* if (!r.push(hash, seq, now))
* ; // full, the caller drops the event
* r.drain([](Slot const& s) { use(s); });
* @endcode
*
* @note Exactly one thread may push and exactly one may drain. Two
* pushers corrupt the ring.
*/
class Ring
{
public:
/**
* Add one event to the ring.
* @param hash Ledger hash to record.
* @param seq Ledger sequence number.
* @param at Time the producer observed the event.
* @return false when the ring is full, in which case nothing was
* stored and the caller must drop the event.
*/
[[nodiscard]] bool
push(uint256 const& hash, LedgerIndex seq, TimePoint at)
{
auto const head = head_.load(std::memory_order_relaxed);
if (head - tail_.load(std::memory_order_acquire) >= kRingCapacity)
return false;
slots_[head & (kRingCapacity - 1)] = Slot{.hash = hash, .seq = seq, .at = at};
head_.store(head + 1, std::memory_order_release);
return true;
}
/**
* Hand every stored event to fn, oldest first, and free their slots.
* @param fn Callable taking Slot const&.
*/
template <class Fn>
void
drain(Fn&& fn)
{
auto tail = tail_.load(std::memory_order_relaxed);
auto const head = head_.load(std::memory_order_acquire);
for (; tail != head; ++tail)
fn(slots_[tail & (kRingCapacity - 1)]);
tail_.store(tail, std::memory_order_release);
}
private:
/**
* Storage, indexed by head_ or tail_ masked to the capacity.
*/
std::array<Slot, kRingCapacity> slots_{};
/**
* Count of events ever pushed. Only the producer writes it.
*/
std::atomic<std::uint64_t> head_{0};
/**
* Count of events ever drained. Only the consumer writes it.
*/
std::atomic<std::uint64_t> tail_{0};
};
/**
* Per-ledger tracking state held in the pending map.
*/
struct LedgerEvent
{
TimePoint recordTime; ///< Time the event was first recorded.
std::uint64_t minute{0}; ///< Minute bucket the event belongs to.
bool weValidated{false}; ///< True if we sent a validation.
bool networkValidated{false}; ///< True if network reached consensus.
bool decided{false}; ///< True once the grace period elapsed.
bool agreed{false}; ///< True if both flags were set.
};
/**
* Counts for one minute of the grid.
*/
struct Bucket
{
std::uint32_t agreed{0}; ///< Agreements decided in this minute.
std::uint32_t total{0}; ///< Events decided in this minute.
};
/**
* Running counts for one rolling window.
*/
struct WindowCount
{
std::uint64_t agreed{0}; ///< Agreements still inside the window.
std::uint64_t total{0}; ///< Events still inside the window.
/**
* Misses still inside the window.
* @return total minus agreed. Derived, so a repair only has to move
* agreed.
*/
[[nodiscard]] std::uint64_t
missed() const
{
return total - agreed;
}
};
/**
* The nine numbers a reader wants, published as one value.
*/
struct Snapshot
{
WindowCount w1h; ///< 1-hour window counts.
WindowCount w24h; ///< 24-hour window counts.
WindowCount w7d; ///< 7-day window counts.
};
/**
* Convert a time point to its minute on the grid.
* @param t Time point to convert.
* @return Whole minutes since the clock's epoch.
*/
static std::uint64_t
minuteOf(TimePoint t);
/**
* Agreement percentage for one window.
* @param w Window counts to divide.
* @return Percentage [0.0, 100.0], or 0.0 when the window is empty.
*/
static double
pct(WindowCount const& w);
/**
* Oldest minute a window of the given length still covers.
* @param minute Newest minute recorded.
* @param span Window length in minutes.
* @return That window's tail minute, floored at zero.
*/
static std::uint64_t
oldestInWindow(std::uint64_t minute, std::size_t span);
/**
* The snapshot readers are currently seeing.
* @return A copy of the published snapshot, so all nine numbers come from
* one reconcile. All zeroes before the first reconcile() publishes.
*/
[[nodiscard]] Snapshot
read() const;
/**
* Put the running counters into a fresh snapshot and publish it.
*/
void
publish();
/**
* Move both rings' contents into pending_.
*/
void
drainRings();
/**
* Fold one drained event into pending_.
* @param s Slot the ring handed over.
* @param ours True if the event came from our own ring.
*/
void
note(Slot const& s, bool ours);
/**
* Decide every pending event past the grace period, repair the ones whose
* other half arrived late, and drop entries too old to repair.
* @param now Current time point.
*/
void
decidePending(TimePoint now);
/**
* Remember a ledger hash as counted, dropping the oldest remembered
* hash once kMaxTalliedEvents is reached.
* @param ledgerHash Hash of the ledger just counted into the totals.
*/
void
noteTallied(uint256 const& ledgerHash);
/**
* Count one decided event in its own bucket and in all three windows.
* @param minute Bucket the event belongs to.
* @param agreed True to count it as an agreement.
*/
void
addToWindows(std::uint64_t minute, bool agreed);
/**
* Turn one already-counted event from a miss into an agreement.
* @param minute Bucket the event was counted in.
*/
void
repairInWindows(std::uint64_t minute);
/**
* Move each window's tail up to the given minute, subtracting whatever
* leaves. The 7-day tail also clears the bucket it passes, because that
* slot is about to be reused.
* @param minute Newest minute to account for.
*/
void
advanceWindows(std::uint64_t minute);
/**
* Walk one window's tail forward, taking each passed bucket back out of
* that window's running counts.
* @param tail The window's tail minute, advanced in place.
* @param target Minute to stop at, the oldest the window still covers.
* @param count The window's running counts to subtract from.
* @param clear True to zero each passed bucket, which only the 7-day
* tail does because it is the tail whose slot gets reused.
*/
void
retireWindow(std::uint64_t& tail, std::uint64_t target, WindowCount& count, bool clear);
/**
* Time source, read on every write and by the reducer.
*/
NowFn now_;
/**
* Events from our own validations. Pushed by the consensus thread.
*/
Ring ourRing_;
/**
* Events from network consensus. Pushed by the ledger master thread.
*/
Ring networkRing_;
/**
* Set while a thread is inside reconcile(). A second caller sees it set
* and returns rather than waiting.
*/
std::atomic_flag reducing_;
/**
* Pending ledger events indexed by ledger hash. Touched only inside
* reconcile(), so it needs no synchronisation.
*/
hash_map<uint256, LedgerEvent> pending_;
/**
* Ledger hashes already counted into the agreement and missed totals.
* Membership survives eviction from pending_, so a ledger reaches the
* totals once. Holds at most kMaxTalliedEvents hashes.
*/
hash_set<uint256> tallied_;
/**
* The hashes in tallied_ in the order they were counted. The front is
* the oldest and is dropped first once the bound is reached.
*/
std::deque<uint256> talliedOrder_;
/**
* One-minute counts, indexed by minute modulo kBuckets7d.
*/
std::array<Bucket, kBuckets7d> buckets_{};
/**
* Running counts for the 1-hour window.
*/
WindowCount c1h_;
/**
* Running counts for the 24-hour window.
*/
WindowCount c24h_;
/**
* Running counts for the 7-day window.
*/
WindowCount c7d_;
/**
* Oldest minute the 1-hour window still counts.
*/
std::uint64_t tail1h_{0};
/**
* Oldest minute the 24-hour window still counts.
*/
std::uint64_t tail24h_{0};
/**
* Oldest minute the 7-day window still counts.
*/
std::uint64_t tail7d_{0};
/**
* Newest minute written to the grid.
*/
std::uint64_t newestMinute_{0};
/**
* False until the first minute is recorded, which is when the tails and
* newestMinute_ get their starting value.
*/
bool started_{false};
/**
* The snapshot readers see. The reducer swaps in a new one each cycle, and
* a reader that took the old one keeps it alive while it reads. Null until
* the first reconcile().
*/
boost::atomic_shared_ptr<Snapshot const> published_;
/**
* Lifetime count of agreements.
*/
std::atomic<std::uint64_t> totalAgreements_{0};
/**
* Lifetime count of misses.
*/
std::atomic<std::uint64_t> totalMissed_{0};
/**
* Agreements at first classification. Backs totalAgreementsEver(); a
* repair never touches it.
*/
std::atomic<std::uint64_t> totalAgreementsGross_{0};
/**
* Misses at first classification. Backs totalMissedEver(); a repair never
* decrements it.
*/
std::atomic<std::uint64_t> totalMissedGross_{0};
/**
* Lifetime count of validations this node sent.
*/
std::atomic<std::uint64_t> totalValidationsSent_{0};
/**
* Lifetime count of network validations observed.
*/
std::atomic<std::uint64_t> totalValidationsChecked_{0};
/**
* Lifetime count of events dropped by a full ring.
*/
std::atomic<std::uint64_t> droppedEvents_{0};
};
} // namespace xrpl::telemetry