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

Brings in the phase-10 revert of the nodestore read-latency histogram plus the
nudb_bytes -> stored_object_bytes rename.

Conflicts in MetricsRegistry.{h,cpp} resolved keeping both intents:

- MetricsRegistry.cpp: dropped everything that existed only to serve the
  reverted nodestore_read_us histogram -- the addSubMillisecondHistogramView()
  helper, its call site, the kSubMillisecondBoundaries array and the
  NodeStoreMetricNames.h include. Kept every view this branch registers
  (consensus round duration, sweep_malloc_trim_us, dns_resolve_latency_ms,
  overlay_dial_latency_ms) and the shared addHistogramView() base helper.
  Took the rename at the storage_detail observe() call site.

- MetricsRegistry.h: took phase-10's move of the four nodestore_state observe
  helpers and their ObserveFn sink from private to public, while keeping this
  branch's enriched Doxygen on observeNodeStoreTotals().

Also corrected the registered-view count in the 09 reference doc: neither side's
arithmetic survives the merge, since this branch adds four views phase-10 never
saw and the revert removes one. Ten views are registered now, not six or seven.
This commit is contained in:
Pratik Mankawde
2026-07-28 16:53:13 +01:00
22 changed files with 1079 additions and 1359 deletions

View File

@@ -541,7 +541,7 @@ public:
env.app().config().getValueFor(SizedItem::TreeCacheAge, std::nullopt)));
}
NodeStoreScheduler scheduler(env.app(), env.app().getJobQueue());
NodeStoreScheduler scheduler(env.app().getJobQueue());
std::string const writableDb = "write";
std::string const archiveDb = "archive";

View File

@@ -4,6 +4,13 @@
#include <test/unit_test/SuiteJournal.h>
#include <xrpld/core/Config.h>
#ifdef XRPL_ENABLE_TELEMETRY
// The four nodestore_state gauge helpers under test, plus the counter type one
// of them reads. Both live in xrpld and are only declared in a
// telemetry-enabled build, so the include is guarded like its uses below.
#include <xrpld/app/ledger/AcquireStats.h>
#include <xrpld/telemetry/MetricsRegistry.h>
#endif
#include <xrpl/basics/Blob.h>
#include <xrpl/basics/ByteUtilities.h>
@@ -21,16 +28,21 @@
#include <xrpl/nodestore/DummyScheduler.h>
#include <xrpl/nodestore/Manager.h>
#include <xrpl/nodestore/NodeObject.h>
#include <xrpl/nodestore/Scheduler.h>
#include <xrpl/nodestore/Types.h>
#include <xrpl/nodestore/detail/DatabaseRotatingImp.h>
#include <xrpl/rdb/DatabaseCon.h>
#include <algorithm>
#include <chrono>
#include <cstddef>
#include <cstdint>
#include <limits>
#include <memory>
#include <optional>
#include <string>
#include <type_traits>
#include <utility>
#include <vector>
namespace xrpl::node_store {
@@ -624,12 +636,6 @@ public:
std::is_same_v<decltype(std::declval<Database const&>().getFetchSize()), std::uint64_t>,
"getFetchSize must be 64-bit");
// A 64-bit counter must be able to represent the byte totals a
// long-lived node reaches. 32 bits cannot.
static_assert(
std::numeric_limits<std::uint64_t>::max() > std::numeric_limits<std::uint32_t>::max(),
"64-bit counters must exceed the 32-bit ceiling");
DummyScheduler scheduler;
beast::TempDir const nodeDb;
@@ -725,11 +731,24 @@ public:
constexpr std::uint64_t kNumStored = 32;
auto const stored = createPredictableBatch(static_cast<int>(kNumStored), 4321);
BEAST_EXPECT(stored.size() == kNumStored);
auto const writesBegan = std::chrono::steady_clock::now();
storeBatch(*db, stored);
auto const wallClockUs =
static_cast<std::uint64_t>(std::chrono::duration_cast<std::chrono::microseconds>(
std::chrono::steady_clock::now() - writesBegan)
.count());
BEAST_EXPECT(db->getStoreCount() == kNumStored);
BEAST_EXPECT(db->getStoreDurationUs() > 0);
// Bounded above by the wall-clock span of the loop that produced it.
// The accumulator only ever sums the backend calls made inside that
// span, so it cannot exceed it. Without this bound the `> 0` above
// would also pass for an accumulator adding a fixed constant per
// insert instead of the measured time.
BEAST_EXPECT(db->getStoreDurationUs() <= wallClockUs);
// Writes must leave the read accumulator alone. This also proves the
// read accessor does not report the write member.
BEAST_EXPECT(db->getFetchDurationUs() == 0);
@@ -765,6 +784,47 @@ public:
//--------------------------------------------------------------------------
/**
* Build a rotating nudb database over two fresh directories.
*
* @param scheduler Scheduler the database keeps a reference to; it must
* outlive the returned database.
* @param writableDir Directory for the writable backend.
* @param archiveDir Directory for the archive backend.
* @return The rotating database, or nullptr if either backend failed.
*/
std::unique_ptr<DatabaseRotatingImp>
makeRotatingDatabase(
Scheduler& scheduler,
beast::TempDir const& writableDir,
beast::TempDir const& archiveDir)
{
Section writableParams;
writableParams.set(Keys::kType, "nudb");
writableParams.set(Keys::kPath, writableDir.path());
Section archiveParams;
archiveParams.set(Keys::kType, "nudb");
archiveParams.set(Keys::kPath, archiveDir.path());
std::shared_ptr<Backend> writableBackend =
Manager::instance().makeBackend(writableParams, megabytes(4), scheduler, journal_);
std::shared_ptr<Backend> archiveBackend =
Manager::instance().makeBackend(archiveParams, megabytes(4), scheduler, journal_);
if (!writableBackend || !archiveBackend)
return nullptr;
writableBackend->open();
archiveBackend->open();
return std::make_unique<DatabaseRotatingImp>(
scheduler,
2,
std::move(writableBackend),
std::move(archiveBackend),
writableParams,
journal_);
}
/**
* Verify the rotating database records its write duration too.
*
@@ -783,35 +843,14 @@ public:
beast::TempDir const writableDir;
beast::TempDir const archiveDir;
Section writableParams;
writableParams.set(Keys::kType, "nudb");
writableParams.set(Keys::kPath, writableDir.path());
Section archiveParams;
archiveParams.set(Keys::kType, "nudb");
archiveParams.set(Keys::kPath, archiveDir.path());
std::shared_ptr<Backend> writableBackend =
Manager::instance().makeBackend(writableParams, megabytes(4), scheduler, journal_);
std::shared_ptr<Backend> archiveBackend =
Manager::instance().makeBackend(archiveParams, megabytes(4), scheduler, journal_);
if (!BEAST_EXPECT(writableBackend) || !BEAST_EXPECT(archiveBackend))
auto rotating = makeRotatingDatabase(scheduler, writableDir, archiveDir);
if (!BEAST_EXPECT(rotating))
return;
writableBackend->open();
archiveBackend->open();
DatabaseRotatingImp rotating(
scheduler,
2,
std::move(writableBackend),
std::move(archiveBackend),
writableParams,
journal_);
// The private fetchNodeObject override hides the public base overload,
// so exercise the rotating store through the Database interface, which
// is also how production callers reach it.
Database& db = rotating;
Database& db = *rotating;
// A fresh rotating database has done no work either.
BEAST_EXPECT(db.getStoreCount() == 0);
@@ -855,6 +894,405 @@ public:
BEAST_EXPECT(db.getStoreDurationUs() == storeDurationAfterWrites);
}
//--------------------------------------------------------------------------
#ifdef XRPL_ENABLE_TELEMETRY
/**
* Recording sink for the nodestore_state gauge helpers.
*
* The helpers take an ObserveFn rather than the OTel observer result
* precisely so a test can hand them somewhere else to put a name and a
* number. Emissions are kept in order, not merged into a map, so a label
* published twice is visible instead of silently overwriting itself.
*/
struct MetricSink
{
/**
* Every (metric label, value) pair published, in emission order.
*/
std::vector<std::pair<std::string, std::int64_t>> emitted;
/**
* Return a sink callable that appends into @ref emitted.
*/
telemetry::MetricsRegistry::ObserveFn
fn()
{
return
[this](char const* name, std::int64_t value) { emitted.emplace_back(name, value); };
}
/**
* Return every published label, sorted, for set comparison.
*/
[[nodiscard]] std::vector<std::string>
names() const
{
std::vector<std::string> out;
out.reserve(emitted.size());
for (auto const& entry : emitted)
out.push_back(entry.first);
std::ranges::sort(out);
return out;
}
/**
* Return the value published for @p name, or std::nullopt when the
* label was not published at all.
*
* Absence and zero must stay distinguishable: a mean over no samples
* is deliberately omitted, while a counter at zero is published.
*
* @param name Metric label to look up.
*/
[[nodiscard]] std::optional<std::int64_t>
value(std::string const& name) const
{
auto const it = std::ranges::find_if(
emitted, [&name](auto const& entry) { return entry.first == name; });
if (it == emitted.end())
return std::nullopt;
return it->second;
}
};
/**
* Build a nudb database with a non-default read-thread count and bundle.
*
* Both are set away from their defaults so the read-queue labels pin real
* forwarding rather than agreeing with a default by accident.
*
* @param dir Directory the store lives in.
* @param scheduler Scheduler the database keeps a reference to; it must
* outlive the returned database.
* @param readThreads Read threads to request.
* @return The database, or nullptr on failure.
*/
std::unique_ptr<Database>
makeMeasuredDatabase(beast::TempDir const& dir, Scheduler& scheduler, int readThreads)
{
Section params;
params.set(Keys::kType, "nudb");
params.set(Keys::kPath, dir.path());
params.set(Keys::kRqBundle, "7");
return Manager::instance().makeDatabase(
megabytes(4), scheduler, readThreads, params, journal_);
}
/**
* Verify the exact `metric` label set observeNodeStoreTotals() publishes,
* and that each derived mean is OMITTED rather than reported as zero
* before there is anything to average.
*
* The 22 label values on this gauge are the change's entire user-visible
* surface. A one-character typo in any of them produces a silently
* disjoint Prometheus series: the metric still exports, the dashboard
* panel goes blank, and nothing else fails.
*/
void
testNodeStoreTotalLabels()
{
testcase("nodestore_state totals labels");
DummyScheduler scheduler;
beast::TempDir const nodeDb;
Section nodeParams;
nodeParams.set(Keys::kType, "nudb");
nodeParams.set(Keys::kPath, nodeDb.path());
std::unique_ptr<Database> db =
Manager::instance().makeDatabase(megabytes(4), scheduler, 2, nodeParams, journal_);
if (!BEAST_EXPECT(db))
return;
// Fresh store: the eight unconditional labels are published, each at
// exactly zero, and NEITHER mean appears.
MetricSink fresh;
telemetry::MetricsRegistry::observeNodeStoreTotals(*db, fresh.fn());
std::vector<std::string> const kFreshLabels{
"node_read_bytes",
"node_reads_duration_us",
"node_reads_hit",
"node_reads_total",
"node_writes",
"node_writes_duration_us",
"node_written_bytes",
"write_load"};
BEAST_EXPECT(fresh.names() == kFreshLabels);
// No label published twice; a map-based sink would have hidden that.
BEAST_EXPECT(fresh.emitted.size() == kFreshLabels.size());
for (auto const& label : kFreshLabels)
BEAST_EXPECT(fresh.value(label) == std::int64_t{0});
// THE omission assertion. A refactor to `mean.value_or(0)` publishes
// these as a believable flat zero on a latency axis; absence is the
// only honest answer before anything has been read or written.
BEAST_EXPECT(!fresh.value("read_mean_us").has_value());
BEAST_EXPECT(!fresh.value("write_mean_us").has_value());
testNodeStoreMeanLabels(*db);
}
/**
* Verify both derived means appear once there is work to average, and
* that each is the quotient of its OWN published numerator and
* denominator.
*
* @param db Database to drive; must be freshly created.
*/
void
testNodeStoreMeanLabels(Database& db)
{
// Reads and writes are given deliberately different counts, so a mean
// wired to the wrong accessor pair -- the adjacent-call-site copy-paste
// risk -- cannot land on the right number.
constexpr std::uint64_t kNumStored = 32;
constexpr std::uint64_t kNumMissing = 8;
auto const stored = createPredictableBatch(static_cast<int>(kNumStored), 24680);
storeBatch(db, stored);
// Fetch each stored object twice, plus a batch of misses, so the read
// count is 2*32 + 8 = 72 against a write count of 32.
for (int pass = 0; pass < 2; ++pass)
{
for (auto const& object : stored)
BEAST_EXPECT(db.fetchNodeObject(object->getHash(), 0) != nullptr);
}
auto const missing = createPredictableBatch(static_cast<int>(kNumMissing), 13579);
for (auto const& object : missing)
BEAST_EXPECT(db.fetchNodeObject(object->getHash(), 0) == nullptr);
MetricSink busy;
telemetry::MetricsRegistry::observeNodeStoreTotals(db, busy.fn());
// Ten labels now: the eight above plus both means.
BEAST_EXPECT(busy.emitted.size() == 10);
BEAST_EXPECT(busy.value("read_mean_us").has_value());
BEAST_EXPECT(busy.value("write_mean_us").has_value());
// Exact counts, so the denominators below are pinned independently.
BEAST_EXPECT(busy.value("node_writes") == std::int64_t{kNumStored});
BEAST_EXPECT(busy.value("node_reads_total") == std::int64_t{2 * kNumStored + kNumMissing});
BEAST_EXPECT(busy.value("node_reads_hit") == std::int64_t{2 * kNumStored});
// The two denominators differ, which is what makes the cross-checks
// below able to catch a swapped accessor pair.
BEAST_EXPECT(busy.value("node_reads_total") != busy.value("node_writes"));
// Each mean is the truncating quotient of two OTHER published labels,
// so this ties the three together without restating the helper's own
// expression. read_mean_us fed the store pair fails it.
auto const quotient = [](std::optional<std::int64_t> total,
std::optional<std::int64_t> count) {
return (total && count && *count != 0) ? std::optional<std::int64_t>{*total / *count}
: std::nullopt;
};
BEAST_EXPECT(
busy.value("read_mean_us") ==
quotient(busy.value("node_reads_duration_us"), busy.value("node_reads_total")));
BEAST_EXPECT(
busy.value("write_mean_us") ==
quotient(busy.value("node_writes_duration_us"), busy.value("node_writes")));
}
/**
* Verify observeWritePathDetail() publishes its four NuDB labels only for
* a measuring backend, and omits the two derived means until an insert
* has happened.
*/
void
testWritePathDetailLabels()
{
testcase("nodestore_state write-path labels");
DummyScheduler scheduler;
// Negative path first: a backend that does not measure its writes must
// publish NOTHING. Zeros here would read as a perfectly idle write
// path on a node whose write path is simply not instrumented.
{
beast::TempDir const memDb;
Section memParams;
memParams.set(Keys::kType, "memory");
memParams.set(Keys::kPath, memDb.path());
std::unique_ptr<Database> mem =
Manager::instance().makeDatabase(megabytes(4), scheduler, 2, memParams, journal_);
if (!BEAST_EXPECT(mem))
return;
auto const batch = createPredictableBatch(8, 111);
storeBatch(*mem, batch);
MetricSink sink;
telemetry::MetricsRegistry::observeWritePathDetail(*mem, sink.fn());
// Cause as well as state: the store really was written to, so the
// emptiness is the std::nullopt branch and not an idle database.
BEAST_EXPECT(sink.emitted.empty());
BEAST_EXPECT(mem->getStoreCount() == 8);
}
beast::TempDir const nodeDb;
Section nodeParams;
nodeParams.set(Keys::kType, "nudb");
nodeParams.set(Keys::kPath, nodeDb.path());
std::unique_ptr<Database> db =
Manager::instance().makeDatabase(megabytes(4), scheduler, 2, nodeParams, journal_);
if (!BEAST_EXPECT(db))
return;
// Measuring backend, nothing written yet: the two instantaneous
// values are published at zero while the two means are omitted. This
// pins the deliberate asymmetry -- zero is meaningful for a gauge and
// meaningless for a mean.
MetricSink fresh;
telemetry::MetricsRegistry::observeWritePathDetail(*db, fresh.fn());
std::vector<std::string> const kFreshLabels{"nudb_insert_max_us", "nudb_writers_in_flight"};
BEAST_EXPECT(fresh.names() == kFreshLabels);
BEAST_EXPECT(fresh.value("nudb_writers_in_flight") == std::int64_t{0});
BEAST_EXPECT(fresh.value("nudb_insert_max_us") == std::int64_t{0});
BEAST_EXPECT(!fresh.value("nudb_insert_mean_us").has_value());
BEAST_EXPECT(!fresh.value("nudb_writer_depth_x100").has_value());
constexpr std::uint64_t kNumStored = 16;
auto const stored = createPredictableBatch(static_cast<int>(kNumStored), 97531);
storeBatch(*db, stored);
MetricSink busy;
telemetry::MetricsRegistry::observeWritePathDetail(*db, busy.fn());
std::vector<std::string> const kBusyLabels{
"nudb_insert_max_us",
"nudb_insert_mean_us",
"nudb_writer_depth_x100",
"nudb_writers_in_flight"};
BEAST_EXPECT(busy.names() == kBusyLabels);
BEAST_EXPECT(busy.emitted.size() == kBusyLabels.size());
// Writers all returned, so the live gauge is back to exactly zero
// while the cumulative maximum stayed up.
BEAST_EXPECT(busy.value("nudb_writers_in_flight") == std::int64_t{0});
// One writing thread, so mean depth is exactly 1.00 and the x100
// fixed-point form must read exactly 100. This is the assertion that
// catches the scale being dropped (it would read 1) or applied twice
// (10000) -- either of which makes the dashboard's divide-by-100 wrong
// by two orders of magnitude.
BEAST_EXPECT(busy.value("nudb_writer_depth_x100") == std::int64_t{100});
}
/**
* Verify the seven acquire_* labels, published unconditionally, and each
* wired to its own counter.
*
* Every expected value below is distinct, so a getter cross-wired to a
* neighbouring counter cannot produce the right number anywhere.
*/
void
testAcquireStatLabels()
{
testcase("nodestore_state acquisition labels");
std::vector<std::string> const kLabels{
"acquire_aborts",
"acquire_aborts_partial",
"acquire_completions",
"acquire_deferrals",
"acquire_give_ups",
"acquire_sweep_evictions",
"acquire_timeouts"};
// A quiet node publishes all seven at zero. Unlike a mean, zero is the
// meaningful "no such event yet" reading for a counter, so omitting
// these would lose the ability to see that nothing happened.
AcquireStats quiet;
MetricSink fresh;
telemetry::MetricsRegistry::observeAcquireStats(quiet, fresh.fn());
BEAST_EXPECT(fresh.names() == kLabels);
BEAST_EXPECT(fresh.emitted.size() == kLabels.size());
for (auto const& label : kLabels)
BEAST_EXPECT(fresh.value(label) == std::int64_t{0});
// Distinct counts per event, and aborts recorded both with and
// without partial work so the subset relationship is exercised: 5
// aborts of which 2 discarded partly built maps.
AcquireStats busy;
for (int i = 0; i < 3; ++i)
busy.recordDeferral();
for (int i = 0; i < 7; ++i)
busy.recordTimeout();
busy.recordGiveUp();
for (int i = 0; i < 2; ++i)
busy.recordAbort(true);
for (int i = 0; i < 3; ++i)
busy.recordAbort(false);
for (int i = 0; i < 11; ++i)
busy.recordCompletion();
for (int i = 0; i < 13; ++i)
busy.recordSweepEviction();
MetricSink sink;
telemetry::MetricsRegistry::observeAcquireStats(busy, sink.fn());
BEAST_EXPECT(sink.names() == kLabels);
BEAST_EXPECT(sink.value("acquire_deferrals") == std::int64_t{3});
BEAST_EXPECT(sink.value("acquire_timeouts") == std::int64_t{7});
BEAST_EXPECT(sink.value("acquire_give_ups") == std::int64_t{1});
BEAST_EXPECT(sink.value("acquire_aborts") == std::int64_t{5});
BEAST_EXPECT(sink.value("acquire_aborts_partial") == std::int64_t{2});
BEAST_EXPECT(sink.value("acquire_completions") == std::int64_t{11});
BEAST_EXPECT(sink.value("acquire_sweep_evictions") == std::int64_t{13});
}
/**
* Verify the four read-queue labels and that the two configured values
* are forwarded rather than defaulted.
*/
void
testReadQueueLabels()
{
testcase("nodestore_state read-queue labels");
DummyScheduler scheduler;
beast::TempDir const nodeDb;
// Three read threads and a bundle of 7, neither of which is the
// default (the bundle default is 4), so a helper reading the wrong
// JSON member cannot agree by coincidence.
auto db = makeMeasuredDatabase(nodeDb, scheduler, 3);
if (!BEAST_EXPECT(db))
return;
MetricSink sink;
telemetry::MetricsRegistry::observeReadQueue(*db, sink.fn());
std::vector<std::string> const kLabels{
"read_queue", "read_request_bundle", "read_threads_running", "read_threads_total"};
BEAST_EXPECT(sink.names() == kLabels);
BEAST_EXPECT(sink.emitted.size() == kLabels.size());
// Nothing has been queued for asynchronous read.
BEAST_EXPECT(sink.value("read_queue") == std::int64_t{0});
// Both configured values, exactly as requested. read_request_bundle is
// the one that would silently read 4 if the helper looked up the wrong
// JSON member, since 4 is the default.
BEAST_EXPECT(sink.value("read_threads_total") == std::int64_t{3});
BEAST_EXPECT(sink.value("read_request_bundle") == std::int64_t{7});
// The running count races with the read threads parking themselves, so
// only its bounds are assertable: present, non-negative, and never
// above the total. Compared as unwrapped values, since comparing two
// optionals would also pass with both absent.
auto const running = sink.value("read_threads_running");
if (BEAST_EXPECT(running.has_value()))
{
BEAST_EXPECT(*running >= 0);
BEAST_EXPECT(*running <= 3);
}
}
#endif // XRPL_ENABLE_TELEMETRY
void
run() override
{
@@ -864,6 +1302,19 @@ public:
testRotatingDurationAccessors();
#ifdef XRPL_ENABLE_TELEMETRY
// The four gauge helpers are declared and defined only in a
// telemetry-enabled build, so their label assertions are guarded the
// same way.
testNodeStoreTotalLabels();
testWritePathDetailLabels();
testAcquireStatLabels();
testReadQueueLabels();
#endif
testConfig();
}
};

View File

@@ -98,22 +98,10 @@ if(telemetry)
HINTS "${opentelemetry-cpp_PACKAGE_FOLDER_RELEASE}/lib"
REQUIRED
)
# The metric side of the in-memory exporter is a SEPARATE archive
# (libopentelemetry_exporter_in_memory_metric.a) with the same
# no-declared-libs problem, so it needs its own find_library. The
# nodestore read-latency histogram tests use it to read exported
# histogram points back and assert per-bucket counts.
find_library(
OTEL_IN_MEMORY_METRIC_EXPORTER_LIB
NAMES opentelemetry_exporter_in_memory_metric
HINTS "${opentelemetry-cpp_PACKAGE_FOLDER_RELEASE}/lib"
REQUIRED
)
target_link_libraries(
xrpl_tests
PRIVATE
"${OTEL_IN_MEMORY_EXPORTER_LIB}"
"${OTEL_IN_MEMORY_METRIC_EXPORTER_LIB}"
opentelemetry-cpp::opentelemetry-cpp
)
# ValidationTracker lives in src/xrpld/ (not libxrpl), so we compile its

View File

@@ -204,4 +204,54 @@ INSTANTIATE_TEST_SUITE_P(
::testing::ValuesIn(backendTypes()),
[](::testing::TestParamInfo<std::string> const& info) { return info.param; });
// The std::nullopt default on the base class, exercised directly on the two
// backends that always exist in every build -- unlike rocksdb, which the
// parameterized suite above only reaches when XRPL_ROCKSDB_AVAILABLE.
//
// Why absence and not zeros: the exporter skips the whole nudb_* label group
// when getWriteStats() is empty (MetricsRegistry.cpp observeWritePathDetail
// returns early). If the base class returned a default-constructed WriteStats
// instead, every non-NuDB node would publish nudb_writers_in_flight=0 and
// nudb_insert_max_us=0 -- a perfectly idle write path, on a node whose write
// path is simply not instrumented. Each assertion below fails against that
// change.
TEST(BackendWriteStats, non_measuring_backends_report_absence_not_zeros)
{
for (auto const& type : {std::string{"memory"}, std::string{"none"}})
{
SCOPED_TRACE("type=" + type);
DummyScheduler scheduler;
beast::Journal const journal{TestSink::instance()};
beast::TempDir const tempDir;
Section params;
params.set("type", type);
params.set("path", tempDir.path());
auto backend = Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
ASSERT_TRUE(backend);
backend->open();
// Absent before any write.
EXPECT_FALSE(backend->getWriteStats().has_value());
// Still absent after real writes. Cause, not just state: the
// backend has genuinely been used, so the absence is the base-class
// default and not an unopened backend.
beast::xor_shift_engine rng(kSeedValue);
auto const batch = createPredictableBatch(16, rng());
storeBatch(*backend, batch);
EXPECT_FALSE(backend->getWriteStats().has_value());
// These backends queue nothing, so their own write load stays 0. The
// pairing matters: absent stats plus a 0 load is what tells the
// exporter "not measured", whereas present stats reading 0 would mean
// "measured, and idle".
EXPECT_EQ(backend->getWriteLoad(), 0);
backend->close();
}
}
} // namespace xrpl::node_store

View File

@@ -253,12 +253,23 @@ TEST_P(NodeStoreDatabaseTest, write_stats_forwarded_from_backend)
ASSERT_TRUE(after.has_value());
EXPECT_EQ(after->insertCount, batch_.size());
// One writer thread, so the depth recorded at each insert is exactly 1.
// This catches a depth accumulator fed the wrong quantity (insertCount's
// running value, or the elapsed microseconds), but it CANNOT catch one fed
// a constant 1, because here the real depth is 1. That case needs genuine
// overlap and is covered by NuDBFactory.cpp's
// write_stats_measure_depth_under_real_overlap.
EXPECT_EQ(after->depthSum, batch_.size());
// No writer is left in flight once the calls have returned.
EXPECT_EQ(after->concurrentWriters, 0u);
// Same live depth through the other accessor on the same object, so the
// two cannot drift onto different fields.
EXPECT_EQ(db->getWriteLoad(), 0);
EXPECT_GT(after->insertTotalUs, 0u);
// A maximum is never below the mean, which fails if the field held the
// minimum or the first sample instead of a running maximum.
// max * n >= sum. Catches a field holding the running MINIMUM, since
// min * n <= sum with equality only when every sample is identical.
// Degenerates when the samples do not vary; the unconditional guarantee
// that the field is a maximum is the non-decreasing check in
// NuDBFactory.cpp's write_stats_accumulate_per_insert.
EXPECT_GE(after->insertMaxUs * after->insertCount, after->insertTotalUs);
}

View File

@@ -4,6 +4,7 @@
#include <xrpl/config/BasicConfig.h>
#include <xrpl/nodestore/DummyScheduler.h>
#include <xrpl/nodestore/Manager.h>
#include <xrpl/nodestore/WriteStats.h>
#include <gtest/gtest.h>
#include <helpers/CaptureSink.h>
@@ -14,7 +15,9 @@
#include <cstddef>
#include <cstdint>
#include <exception>
#include <latch>
#include <memory>
#include <optional>
#include <string>
#include <thread>
#include <utility>
@@ -56,6 +59,48 @@ runRoundTrip(Section const& params, std::size_t expectedBlocksize)
EXPECT_EQ(batch, copy);
}
/**
* Threads used by the overlapping-insert round below.
*/
constexpr std::uint64_t kOverlapThreads = 8;
/**
* Inserts each of those threads performs per round.
*/
constexpr std::uint64_t kOverlapPerThread = 50;
/**
* Run one round of deliberately overlapping inserts against @p backend.
*
* All kOverlapThreads threads are released from a single latch, so they reach
* doInsert() together rather than one after another; staggered starts are what
* would let every insert run end to end and never overlap.
*
* @param backend Backend to insert into. Must be open.
* @param round Round index, mixed into the seeds so every round writes
* fresh keys and no insert takes the duplicate short-circuit.
*/
void
runOverlappingInsertRound(Backend& backend, int round)
{
std::latch start(static_cast<std::ptrdiff_t>(kOverlapThreads));
std::vector<std::thread> threads;
threads.reserve(kOverlapThreads);
for (auto t = 0uz; t < kOverlapThreads; ++t)
{
threads.emplace_back([&backend, &start, t, round] {
auto const batch = createPredictableBatch(
kOverlapPerThread, 1000 + t + (static_cast<std::uint64_t>(round) * 100'000));
start.arrive_and_wait();
for (auto const& obj : batch)
backend.store(obj);
});
}
for (auto& th : threads)
th.join();
}
} // namespace
TEST(NuDBFactory, default_block_size)
@@ -290,6 +335,14 @@ TEST(NuDBFactory, write_stats_accumulate_per_insert)
// Exactly 10 inserts must be counted as 10, and depthSum must be 10
// because a single-threaded caller is always the only writer, so the
// depth recorded at each insert is exactly 1.
//
// What this pins and what it cannot: on one thread depthSum == insertCount
// catches an accumulator fed the wrong quantity -- fed insertCount it
// would read 1+2+...+n, and fed the elapsed time it would read the
// microseconds. It does NOT catch depthSum being fed a constant 1, which
// is indistinguishable here because the real depth IS 1. That bug is the
// one that would silently zero every derived wait time, and it is caught
// by write_stats_measure_depth_under_real_overlap below.
constexpr std::uint64_t kFirstBatch = 10;
auto const batch = createPredictableBatch(kFirstBatch, 12345);
storeBatch(*backend, batch);
@@ -301,11 +354,12 @@ TEST(NuDBFactory, write_stats_accumulate_per_insert)
EXPECT_EQ(after->depthSum, kFirstBatch);
EXPECT_GT(after->insertTotalUs, 0u);
EXPECT_GT(after->insertMaxUs, 0u);
// The largest single insert cannot exceed the sum of all of them.
EXPECT_LE(after->insertMaxUs, after->insertTotalUs);
// A maximum is never below the mean. This fails if the field were
// holding the minimum, or the first or last sample, instead of the
// running maximum.
// A maximum is never below the mean, so max * n >= sum. Catches a field
// fed the running MINIMUM: min * n <= sum, with equality only when every
// sample is identical, so any variation at all makes the two orderings
// exclusive. It does not discriminate when the samples happen not to
// vary; the running-maximum property is pinned unconditionally by the
// non-decreasing check after the second batch below.
EXPECT_GE(after->insertMaxUs * after->insertCount, after->insertTotalUs);
// No writer remains in flight once the calls have returned.
EXPECT_EQ(after->concurrentWriters, 0u);
@@ -334,6 +388,22 @@ TEST(NuDBFactory, write_stats_accumulate_per_insert)
// Negative path: NuDB reports key_exists for a duplicate key and doInsert
// deliberately does not treat that as an error. The accounting must still
// run, and in particular the depth must come back down.
//
// What this does NOT cover, stated plainly because a comment claiming absent
// coverage is worse than none: this is not the throwing path. nudb::insert()
// sets error::key_exists and RETURNS (nudb/impl/basic_store.ipp:294, :307,
// :329), and doInsert() filters exactly that code out before it would throw
// (NuDBFactory.cpp:283), so a duplicate key takes the identical non-throwing
// control flow as a fresh key. It reaches the ScopeExit guard by the same
// route the happy path does.
//
// The throwing path -- where the guard is the only reason the depth comes
// back down -- is not reachable from a unit test: it needs nudb::insert() to
// fail with something other than key_exists (an I/O or allocation failure
// inside the library), which cannot be induced through the Backend interface
// without a fault-injection seam that does not exist. Its RAII contract is
// covered generically instead: src/tests/libxrpl/basics/scope.cpp:34-45 proves
// ScopeExit runs its function during unwinding.
TEST(NuDBFactory, write_stats_count_duplicate_key_inserts)
{
beast::TempDir const tempDir;
@@ -353,6 +423,7 @@ TEST(NuDBFactory, write_stats_count_duplicate_key_inserts)
if (!first.has_value())
FAIL() << "nudb must report write stats";
ASSERT_EQ(first->insertCount, kBatchSize);
ASSERT_EQ(first->depthSum, kBatchSize);
// Re-storing the identical batch writes nothing new, but each call is
// still an insert attempt that entered and left the backend.
@@ -361,16 +432,40 @@ TEST(NuDBFactory, write_stats_count_duplicate_key_inserts)
auto const second = backend->getWriteStats();
if (!second.has_value())
FAIL() << "nudb must report write stats after re-storing";
EXPECT_EQ(second->insertCount, kBatchSize * 2);
EXPECT_EQ(second->depthSum, kBatchSize * 2);
// The depth returned to zero, so the early-return error path did not
// The duplicate round is counted, so a key_exists early return is not
// skipping the accounting. Written as the first snapshot plus the batch
// size rather than as one product, because the two sides must differ by
// exactly the second round: an implementation that counted only the
// rounds that stored new data would leave these equal.
EXPECT_EQ(second->insertCount, first->insertCount + kBatchSize);
EXPECT_EQ(second->depthSum, first->depthSum + kBatchSize);
// The depth returned to zero, so the key_exists early return did not
// leak a writer.
EXPECT_EQ(second->concurrentWriters, 0u);
EXPECT_EQ(backend->getWriteLoad(), 0);
backend->close();
}
TEST(NuDBFactory, write_stats_observe_concurrent_writers)
// depthSum is the L in Little's Law: mean depth L and mean insert time W give
// service time S = W / L, and the queuing time the whole diagnosis rests on is
// W - S. If depthSum were fed a constant 1 instead of the observed depth then
// L would read exactly 1.0, S would equal W, and every derived wait would read
// 0 -- a stalled write path indistinguishable from a healthy one, with nothing
// on any dashboard looking wrong.
//
// A single-threaded test cannot see that bug, because there the real depth IS
// 1. This test forces genuine overlap so the correct implementation records a
// depth above 1 and the constant-1 implementation cannot.
//
// Why the overlap is reachable and not merely hoped for: NuDB takes one global
// mutex for the entire insert, and doInsert() reads the depth BEFORE entering
// it. So while one thread is inside an insert, every other thread that reaches
// doInsert() records a depth of at least 2 and then blocks. All threads are
// released from one latch, and the round is retried until the overlap is
// observed -- so a constant-1 implementation exhausts every round and fails,
// while the real one satisfies it as soon as any two inserts overlap.
TEST(NuDBFactory, write_stats_measure_depth_under_real_overlap)
{
beast::TempDir const tempDir;
auto const params = makeSection(tempDir.path());
@@ -381,36 +476,57 @@ TEST(NuDBFactory, write_stats_observe_concurrent_writers)
ASSERT_TRUE(backend);
backend->open();
// Four threads insert distinct objects concurrently. The exact peak
// depth is racy, but two invariants are not: every insert is counted,
// and depthSum is at least insertCount because depth is >= 1 per
// insert.
constexpr std::uint64_t kThreads = 4;
constexpr std::uint64_t kPerThread = 50;
// Bounded so a genuine regression fails instead of hanging. Each round
// runs kOverlapThreads * kOverlapPerThread inserts through one global
// mutex, so one round already gives the correct implementation many
// chances to overlap.
constexpr int kMaxRounds = 20;
std::vector<std::thread> threads;
threads.reserve(kThreads);
for (auto t = 0uz; t < kThreads; ++t)
std::uint64_t completedRounds = 0;
std::optional<WriteStats> stats;
for (auto round = 0; round < kMaxRounds; ++round)
{
threads.emplace_back([&backend, t] {
auto const batch = createPredictableBatch(kPerThread, 1000 + t);
for (auto const& obj : batch)
backend->store(obj);
});
}
for (auto& th : threads)
th.join();
runOverlappingInsertRound(*backend, round);
++completedRounds;
stats = backend->getWriteStats();
if (!stats.has_value())
FAIL() << "nudb must report write stats after concurrent inserts";
if (stats->depthSum > stats->insertCount)
break;
}
auto const stats = backend->getWriteStats();
if (!stats.has_value())
FAIL() << "nudb must report write stats after concurrent inserts";
EXPECT_EQ(stats->insertCount, kThreads * kPerThread);
EXPECT_GE(stats->depthSum, stats->insertCount);
FAIL() << "no round produced write stats";
// Every insert of every round is counted exactly once. A lost increment
// under contention fails this.
EXPECT_EQ(stats->insertCount, completedRounds * kOverlapThreads * kOverlapPerThread);
// THE assertion this test exists for: strictly greater, so a depthSum fed
// a constant 1 (or fed nothing, or fed insertCount's own delta) cannot
// satisfy it however many rounds run.
EXPECT_GT(stats->depthSum, stats->insertCount)
<< "depthSum must record the observed depth, not a constant 1; rounds run="
<< completedRounds;
// Upper bound with teeth: at most kThreads writers can be inside an
// insert at once, so no single insert can observe a depth above kThreads.
// A missing fetch_sub in recordInsert() would let the gauge climb once
// per insert, giving a depthSum near insertCount squared over two --
// vastly over this bound at these counts.
EXPECT_LE(stats->depthSum, stats->insertCount * kOverlapThreads);
// State plus cause: the gauge is back to exactly zero, so every one of
// the increments taken above was matched by its decrement. Exactly 0 and
// not "small": a single leaked writer strands getWriteLoad() nonzero for
// the life of the process, which gates history acquisition.
EXPECT_EQ(stats->concurrentWriters, 0u);
EXPECT_EQ(backend->getWriteLoad(), 0);
EXPECT_GT(stats->insertMaxUs, 0u);
// Depth cannot exceed the number of threads that could be inside the
// insert at once, so the mean depth is bounded by kThreads.
EXPECT_LE(stats->depthSum, stats->insertCount * kThreads);
backend->close();
}
@@ -431,15 +547,29 @@ TEST(NuDBFactory, write_load_reports_writer_depth)
// After writes complete the depth returns to 0 rather than staying
// elevated, because this is an instantaneous gauge and not a counter.
// Exactly 0 and not merely small: were getWriteLoad() to return one of
// the cumulative fields instead of the live depth -- insertCount would
// read 5 here, insertTotalUs some microsecond total -- this fails.
auto const batch = createPredictableBatch(5, 777);
storeBatch(*backend, batch);
EXPECT_EQ(backend->getWriteLoad(), 0);
// The value must stay far below the history-acquisition cutoff that
// LedgerMaster applies (kMaxWriteLoadAcquire), or history acquisition
// would silently stop. Depth is bounded by the writing threads.
constexpr int kMaxWriteLoadAcquire = 8192;
EXPECT_LT(backend->getWriteLoad(), kMaxWriteLoadAcquire);
// Same value as the write-stats snapshot reports, since both read the one
// depth atomic. Catches the two accessors drifting onto different fields.
auto const stats = backend->getWriteStats();
if (!stats.has_value())
FAIL() << "nudb must report write stats";
EXPECT_EQ(static_cast<std::uint64_t>(backend->getWriteLoad()), stats->concurrentWriters);
// The cumulative fields did move, so the 0 above is the gauge being
// instantaneous and not the backend having done nothing.
EXPECT_EQ(stats->insertCount, 5u);
// NOTE. LedgerMaster gates history acquisition on getWriteLoad() staying
// below kMaxWriteLoadAcquire (8192), declared static constexpr inside
// src/xrpld/app/ledger/detail/LedgerMaster.cpp and so unreachable from
// this binary. Depth is bounded by the number of writing threads, which
// cannot approach that figure, so the coupling is recorded here rather
// than asserted against a literal copy of the constant that could drift.
backend->close();
}

View File

@@ -446,7 +446,8 @@ TEST(MetricsRegistryScaledMean, default_scale_is_one)
{
// The two-argument form is the latency case and must not scale silently;
// if the default were 100 every published latency would be 100x wrong.
EXPECT_EQ(Registry::scaledMean(360, 8), Registry::scaledMean(360, 8, 1));
// 360/8 is 45 by hand -- an independent literal, not a restatement of the
// implementation. A default of 100 would read 4500 here.
EXPECT_EQ(Registry::scaledMean(360, 8), 45);
}

View File

@@ -1,526 +0,0 @@
/**
* @file NodeStoreMetricNames.cpp
* Unit tests for the nodestore read-latency histogram wiring.
*
* Two independent groups, split by what they can link:
*
* 1. The shared name/label constants and the three record-site helpers from
* `<xrpl/telemetry/NodeStoreMetricNames.h>`. Header-only and free of any
* OTel dependency, so these run in **both** builds. That matters: the
* constants are what keeps the bucket-view registration in
* MetricsRegistry.cpp and the record site in NodeStoreScheduler.cpp
* agreeing on one instrument name, and a divergence there silently drops
* the sub-millisecond bucket override.
*
* 2. An end-to-end record-and-read-back over a real SDK MeterProvider fitted
* with the same explicit sub-millisecond boundaries production registers,
* asserting the exact bucket counts a set of known latencies must land in.
* Guarded on XRPL_ENABLE_TELEMETRY because the metrics SDK headers only
* exist in that build.
*
* Why the second group does not drive NodeStoreScheduler directly: that class
* lives in xrpld (`src/xrpld/app/main/`) and its onFetch() needs a live
* JobQueue plus a ServiceRegistry, neither of which the standalone xrpl_tests
* binary can supply -- the same reason MetricsRegistry.cpp is only compiled
* into this binary on the no-op path (see src/tests/libxrpl/CMakeLists.txt).
* What is testable here is everything that decides *what* gets recorded: the
* instrument name, the two label values, the negative-value guard, and the
* bucket ladder the value is filed into. The remaining step -- that
* Database::fetchNodeObject actually reaches onFetch -- is covered by the
* existing nodestore suites, which already exercise that call path.
*/
#include <xrpl/telemetry/NodeStoreMetricNames.h>
#include <gtest/gtest.h>
#include <string_view>
namespace {
using namespace xrpl::telemetry;
// ---------------------------------------------------------------------------
// Group 1: shared constants and record-site helpers. Compile-time first, so a
// regression is a build failure rather than only a test failure; the runtime
// duplicates below name the offending case when one does fail.
// ---------------------------------------------------------------------------
// The instrument name is the contract between the two call sites. Pinned to
// the exact literal: bare lower snake_case, no `xrpld_` prefix (no metric in
// this codebase carries one), and the `_us` suffix stating the unit.
static_assert(std::string_view{kNodeStoreReadUs} == "nodestore_read_us");
// Label keys. `found` rather than `was_found` and `fetch_type` rather than
// `type`, matching what the dashboards query.
static_assert(std::string_view{kFetchTypeLabel} == "fetch_type");
static_assert(std::string_view{kFetchFoundLabel} == "found");
// Label values.
static_assert(std::string_view{kFetchTypeAsync} == "async");
static_assert(std::string_view{kFetchTypeSync} == "sync");
static_assert(std::string_view{kFetchFoundTrue} == "true");
static_assert(std::string_view{kFetchFoundFalse} == "false");
// The helpers map each input to exactly one value, and the two arms differ.
static_assert(std::string_view{fetchTypeLabelValue(true)} == "async");
static_assert(std::string_view{fetchTypeLabelValue(false)} == "sync");
static_assert(std::string_view{fetchFoundLabelValue(true)} == "true");
static_assert(std::string_view{fetchFoundLabelValue(false)} == "false");
// The negative guard. Zero is admitted on purpose -- a page-cache hit really
// can round to 0 us -- while anything below it is refused.
static_assert(shouldRecordFetchLatency(0));
static_assert(shouldRecordFetchLatency(1));
static_assert(shouldRecordFetchLatency(25'000));
static_assert(!shouldRecordFetchLatency(-1));
} // namespace
TEST(NodeStoreMetricNames, instrument_name_is_the_exact_shared_literal)
{
// Both the view registration (MetricsRegistry.cpp) and the record site
// (NodeStoreScheduler.cpp) read this one constant. If it changes, the
// dashboard query and the reference doc must change with it, so the exact
// string is asserted rather than merely its shape.
EXPECT_EQ(std::string_view{kNodeStoreReadUs}, "nodestore_read_us");
EXPECT_EQ(std::string_view{kNodeStoreReadUs}.size(), 17u);
// No `xrpld_` prefix: verified against the live metric surface, where 0 of
// 537 exported names carry one. A prefix here would make this the only
// odd metric out and break every dashboard that globs the family.
EXPECT_FALSE(std::string_view{kNodeStoreReadUs}.starts_with("xrpld_"));
// The `_us` suffix is load-bearing: FetchReport::elapsed is
// std::chrono::microseconds, and a name implying milliseconds would make
// every reading 1000x wrong to a reader.
EXPECT_TRUE(std::string_view{kNodeStoreReadUs}.ends_with("_us"));
// The description must name the unit too, since that is all a Prometheus
// consumer sees alongside the metric.
EXPECT_EQ(
std::string_view{kNodeStoreReadUsDesc}, "NodeStore backend fetch latency in microseconds");
}
TEST(NodeStoreMetricNames, label_keys_and_values_are_the_exact_literals)
{
EXPECT_EQ(std::string_view{kFetchTypeLabel}, "fetch_type");
EXPECT_EQ(std::string_view{kFetchFoundLabel}, "found");
EXPECT_EQ(std::string_view{kFetchTypeAsync}, "async");
EXPECT_EQ(std::string_view{kFetchTypeSync}, "sync");
EXPECT_EQ(std::string_view{kFetchFoundTrue}, "true");
EXPECT_EQ(std::string_view{kFetchFoundFalse}, "false");
// The two keys must differ, or one label would overwrite the other in the
// attribute map and a whole dimension would vanish.
EXPECT_NE(std::string_view{kFetchTypeLabel}, std::string_view{kFetchFoundLabel});
}
TEST(NodeStoreMetricNames, helpers_map_each_input_to_its_own_value)
{
// Positive path for both arms of both helpers.
EXPECT_EQ(std::string_view{fetchTypeLabelValue(true)}, "async");
EXPECT_EQ(std::string_view{fetchTypeLabelValue(false)}, "sync");
EXPECT_EQ(std::string_view{fetchFoundLabelValue(true)}, "true");
EXPECT_EQ(std::string_view{fetchFoundLabelValue(false)}, "false");
// Cause, not just state: the two arms are genuinely distinct, so a
// copy-paste that returned the same value for both would fail here rather
// than quietly collapsing async and sync into one series.
EXPECT_NE(
std::string_view{fetchTypeLabelValue(true)}, std::string_view{fetchTypeLabelValue(false)});
EXPECT_NE(
std::string_view{fetchFoundLabelValue(true)},
std::string_view{fetchFoundLabelValue(false)});
// Each helper returns one of its own two constants and never the other
// helper's, which is what keeps the two dimensions independent.
EXPECT_EQ(fetchTypeLabelValue(true), kFetchTypeAsync);
EXPECT_EQ(fetchTypeLabelValue(false), kFetchTypeSync);
EXPECT_EQ(fetchFoundLabelValue(true), kFetchFoundTrue);
EXPECT_EQ(fetchFoundLabelValue(false), kFetchFoundFalse);
}
TEST(NodeStoreMetricNames, latency_guard_admits_zero_and_refuses_negatives)
{
// Negative path -- the reason the guard exists. The OTel SDK drops a
// negative histogram value AND logs a warning for it; on a per-fetch path
// that is a log flood, so the sample is filtered before it gets there.
EXPECT_FALSE(shouldRecordFetchLatency(-1));
EXPECT_FALSE(shouldRecordFetchLatency(-1'000));
// Zero must NOT be filtered: a read served from the page cache genuinely
// truncates to 0 us, and suppressing it would hide the fastest reads and
// bias the whole distribution upward.
EXPECT_TRUE(shouldRecordFetchLatency(0));
// Ordinary and cold-tail values pass.
EXPECT_TRUE(shouldRecordFetchLatency(1));
EXPECT_TRUE(shouldRecordFetchLatency(9));
EXPECT_TRUE(shouldRecordFetchLatency(250));
EXPECT_TRUE(shouldRecordFetchLatency(30'000));
// The boundary is exactly at zero, not near it.
EXPECT_TRUE(shouldRecordFetchLatency(0));
EXPECT_FALSE(shouldRecordFetchLatency(-1));
}
// ---------------------------------------------------------------------------
// Group 2: record into a real histogram carrying production's explicit
// sub-millisecond boundaries, then read the exported point back and assert the
// exact per-bucket counts.
//
// This is what proves the signal is usable rather than merely emitted. The
// SDK's default boundaries begin at 0/5/10/25... but top out at 10,000, and
// the microsecond ladder used by the other duration histograms begins at 100
// us -- above the entire range a warm read occupies. Under that ladder every
// warm read files into bucket 0 and the distribution reads flat. The
// assertions below pin warm reads into distinct low buckets, which is exactly
// the property the sub-millisecond ladder exists to provide.
// ---------------------------------------------------------------------------
#ifdef XRPL_ENABLE_TELEMETRY
#include <opentelemetry/context/context.h>
#include <opentelemetry/exporters/memory/in_memory_metric_data.h>
#include <opentelemetry/exporters/memory/in_memory_metric_exporter_factory.h>
#include <opentelemetry/metrics/meter.h>
#include <opentelemetry/nostd/unique_ptr.h>
#include <opentelemetry/nostd/variant.h>
#include <opentelemetry/sdk/metrics/aggregation/aggregation_config.h>
#include <opentelemetry/sdk/metrics/data/point_data.h>
#include <opentelemetry/sdk/metrics/export/periodic_exporting_metric_reader_factory.h>
#include <opentelemetry/sdk/metrics/export/periodic_exporting_metric_reader_options.h>
#include <opentelemetry/sdk/metrics/instruments.h>
#include <opentelemetry/sdk/metrics/meter_provider.h>
#include <opentelemetry/sdk/metrics/meter_provider_factory.h>
#include <opentelemetry/sdk/metrics/view/instrument_selector_factory.h>
#include <opentelemetry/sdk/metrics/view/meter_selector_factory.h>
#include <opentelemetry/sdk/metrics/view/view_factory.h>
#include <opentelemetry/sdk/metrics/view/view_registry.h>
#include <array>
#include <chrono>
#include <cstddef>
#include <memory>
#include <string>
#include <utility>
#include <vector>
namespace {
namespace metric_sdk = opentelemetry::sdk::metrics;
namespace in_memory = opentelemetry::exporter::memory;
/**
* The same edges MetricsRegistry.cpp's kSubMillisecondBoundaries holds.
*
* Deliberately a second, independent copy rather than an include of the
* production array: that constant lives in an unnamed namespace inside
* MetricsRegistry.cpp and is unreachable from here, and re-deriving the edges
* from the implementation would make the bucket-index assertions below
* tautological. Written out by hand, they pin the ladder -- so silently
* re-tuning an edge in production without revisiting this file fails here.
*/
constexpr std::array kExpectedBoundaries{
1.0,
2.0,
5.0,
10.0,
25.0,
50.0,
100.0,
250.0,
500.0,
1'000.0,
5'000.0,
25'000.0};
/**
* Meter identity used by the production view selector, so the view this
* fixture registers matches the instrument the fixture creates.
*/
constexpr char kMeterName[] = "xrpld";
constexpr char kMeterVersion[] = "1.0.0";
/**
* A MeterProvider carrying one explicit-bucket view for kNodeStoreReadUs and
* an in-memory exporter, so a test can record values and read the resulting
* histogram point back without any network or OTLP involvement.
*
* Mirrors what MetricsRegistry::initExporterAndProvider() builds for this
* instrument, minus the OTLP exporter.
*/
class HistogramFixture
{
public:
HistogramFixture()
{
// The view: same instrument type, same name, same meter selector and
// the same boundaries production registers.
auto config = std::make_shared<metric_sdk::HistogramAggregationConfig>();
config->boundaries_ = {kExpectedBoundaries.begin(), kExpectedBoundaries.end()};
auto views = std::make_unique<metric_sdk::ViewRegistry>();
views->AddView(
metric_sdk::InstrumentSelectorFactory::Create(
metric_sdk::InstrumentType::kHistogram, kNodeStoreReadUs, ""),
metric_sdk::MeterSelectorFactory::Create(kMeterName, kMeterVersion, ""),
metric_sdk::ViewFactory::Create(
kNodeStoreReadUs, "", metric_sdk::AggregationType::kHistogram, config));
provider_ = metric_sdk::MeterProviderFactory::Create(std::move(views));
// A long export interval keeps the background thread from exporting
// on its own schedule; the test drives collection via ForceFlush().
metric_sdk::PeriodicExportingMetricReaderOptions readerOpts;
readerOpts.export_interval_millis = std::chrono::milliseconds(600'000);
readerOpts.export_timeout_millis = std::chrono::milliseconds(5'000);
provider_->AddMetricReader(
metric_sdk::PeriodicExportingMetricReaderFactory::Create(
in_memory::InMemoryMetricExporterFactory::Create(data_), readerOpts));
histogram_ = provider_->GetMeter(kMeterName, kMeterVersion)
->CreateDoubleHistogram(kNodeStoreReadUs, kNodeStoreReadUsDesc);
}
~HistogramFixture()
{
provider_->Shutdown();
}
HistogramFixture(HistogramFixture const&) = delete;
HistogramFixture&
operator=(HistogramFixture const&) = delete;
/**
* Record one latency with the exact label set the production record site
* attaches, built through the same two helpers.
*
* @param elapsedUs Latency in microseconds.
* @param isAsync True for an async (prefetch) read.
* @param wasFound True when the object was found.
*/
void
record(double elapsedUs, bool isAsync, bool wasFound)
{
histogram_->Record(
elapsedUs,
{{kFetchTypeLabel, std::string(fetchTypeLabelValue(isAsync))},
{kFetchFoundLabel, std::string(fetchFoundLabelValue(wasFound))}},
opentelemetry::context::Context{});
}
/**
* Flush the reader, then return every exported histogram point for
* kNodeStoreReadUs keyed by its attribute set.
*/
[[nodiscard]] in_memory::SimpleAggregateInMemoryMetricData::AttributeToPoint const&
collect()
{
provider_->ForceFlush();
return data_->Get(kMeterName, kNodeStoreReadUs);
}
private:
/**
* Sink the in-memory exporter writes each collection into.
*/
std::shared_ptr<in_memory::SimpleAggregateInMemoryMetricData> data_ =
std::make_shared<in_memory::SimpleAggregateInMemoryMetricData>();
/**
* Provider owning the view registry and the in-memory reader.
*/
std::shared_ptr<metric_sdk::MeterProvider> provider_;
/**
* The instrument under test.
*/
opentelemetry::nostd::unique_ptr<opentelemetry::metrics::Histogram<double>> histogram_;
};
/**
* Extract the HistogramPointData for the one series matching @p isAsync and
* @p wasFound, or nullptr when no such series was exported.
*
* @param points Exported points keyed by attribute set.
* @param isAsync fetch_type dimension to match.
* @param wasFound found dimension to match.
*/
[[nodiscard]] metric_sdk::HistogramPointData const*
findPoint(
in_memory::SimpleAggregateInMemoryMetricData::AttributeToPoint const& points,
bool isAsync,
bool wasFound)
{
for (auto const& [attributes, point] : points)
{
auto const type = attributes.find(kFetchTypeLabel);
auto const found = attributes.find(kFetchFoundLabel);
if (type == attributes.end() || found == attributes.end())
continue;
if (opentelemetry::nostd::get<std::string>(type->second) != fetchTypeLabelValue(isAsync) ||
opentelemetry::nostd::get<std::string>(found->second) != fetchFoundLabelValue(wasFound))
continue;
return &opentelemetry::nostd::get<metric_sdk::HistogramPointData>(point);
}
return nullptr;
}
} // namespace
TEST(NodeStoreReadHistogram, view_applies_the_sub_millisecond_boundaries)
{
HistogramFixture fixture;
fixture.record(9.0, /*isAsync=*/false, /*wasFound=*/true);
auto const* point = findPoint(fixture.collect(), /*isAsync=*/false, /*wasFound=*/true);
ASSERT_NE(point, nullptr);
// The exported point must carry OUR boundaries, not the SDK defaults. This
// is the assertion that catches a name mismatch between the view selector
// and the instrument: on a mismatch the view is never applied and the
// default ladder appears here instead.
ASSERT_EQ(point->boundaries_.size(), kExpectedBoundaries.size());
for (std::size_t i = 0; i < kExpectedBoundaries.size(); ++i)
EXPECT_EQ(point->boundaries_[i], kExpectedBoundaries[i]) << "boundary index " << i;
// The first edge is 1 us, three orders of magnitude below the microsecond
// ladder's 100 us first edge. That difference is the entire point: without
// it a warm read cannot be distinguished from an instant one.
EXPECT_EQ(point->boundaries_.front(), 1.0);
EXPECT_EQ(point->boundaries_.back(), 25'000.0);
}
TEST(NodeStoreReadHistogram, warm_reads_land_in_distinct_low_buckets)
{
HistogramFixture fixture;
// Four warm latencies, each chosen to fall in a different low bucket.
// Bucket i counts values in (boundaries[i-1], boundaries[i]].
// 0.5 us -> bucket 0 ( <= 1 )
// 3 us -> bucket 2 ( 2 < v <= 5 )
// 9 us -> bucket 3 ( 5 < v <= 10 )
// 40 us -> bucket 5 ( 25 < v <= 50 )
fixture.record(0.5, /*isAsync=*/false, /*wasFound=*/true);
fixture.record(3.0, /*isAsync=*/false, /*wasFound=*/true);
fixture.record(9.0, /*isAsync=*/false, /*wasFound=*/true);
fixture.record(40.0, /*isAsync=*/false, /*wasFound=*/true);
auto const* point = findPoint(fixture.collect(), /*isAsync=*/false, /*wasFound=*/true);
ASSERT_NE(point, nullptr);
// Exact counts, bucket by bucket -- not merely "the total is 4". Under the
// microsecond ladder all four would sit in bucket 0 and this test would
// fail, which is precisely the regression it guards against.
ASSERT_EQ(point->counts_.size(), kExpectedBoundaries.size() + 1);
EXPECT_EQ(point->counts_[0], 1u); // 0.5 us
EXPECT_EQ(point->counts_[1], 0u);
EXPECT_EQ(point->counts_[2], 1u); // 3 us
EXPECT_EQ(point->counts_[3], 1u); // 9 us
EXPECT_EQ(point->counts_[4], 0u);
EXPECT_EQ(point->counts_[5], 1u); // 40 us
for (std::size_t i = 6; i < point->counts_.size(); ++i)
EXPECT_EQ(point->counts_[i], 0u) << "bucket " << i << " should be empty";
// Aggregate state must agree with the per-bucket detail.
EXPECT_EQ(point->count_, 4u);
EXPECT_EQ(opentelemetry::nostd::get<double>(point->sum_), 52.5);
EXPECT_EQ(opentelemetry::nostd::get<double>(point->min_), 0.5);
EXPECT_EQ(opentelemetry::nostd::get<double>(point->max_), 40.0);
}
TEST(NodeStoreReadHistogram, a_cold_read_lands_in_the_tail_not_the_ceiling)
{
HistogramFixture fixture;
// A cold read and an outlier beyond the top edge. 800 us falls in bucket 9
// ( 500 < v <= 1000 ); 30000 us exceeds the 25000 top edge and so lands in
// the overflow bucket, index 12.
fixture.record(800.0, /*isAsync=*/false, /*wasFound=*/true);
fixture.record(30'000.0, /*isAsync=*/false, /*wasFound=*/true);
auto const* point = findPoint(fixture.collect(), /*isAsync=*/false, /*wasFound=*/true);
ASSERT_NE(point, nullptr);
EXPECT_EQ(point->counts_[9], 1u); // 800 us -- resolved, not saturated
EXPECT_EQ(point->counts_[12], 1u); // 30 ms -- overflow bucket
EXPECT_EQ(point->count_, 2u);
EXPECT_EQ(opentelemetry::nostd::get<double>(point->max_), 30'000.0);
}
TEST(NodeStoreReadHistogram, the_two_labels_split_the_series_four_ways)
{
HistogramFixture fixture;
// One record per (fetch_type, found) combination, each with a distinct
// latency so the series cannot be confused with one another.
fixture.record(3.0, /*isAsync=*/true, /*wasFound=*/true);
fixture.record(9.0, /*isAsync=*/true, /*wasFound=*/false);
fixture.record(40.0, /*isAsync=*/false, /*wasFound=*/true);
fixture.record(800.0, /*isAsync=*/false, /*wasFound=*/false);
auto const& points = fixture.collect();
// Four distinct label sets means four distinct time series. If either
// label were dropped or misspelled these would collapse into fewer.
EXPECT_EQ(points.size(), 4u);
struct Expected
{
bool isAsync;
bool wasFound;
double value;
std::size_t bucket;
};
// Each series holds exactly its own one sample, in its own bucket. This is
// the assertion that a label mix-up would break: swapping two values would
// put a sample in the wrong series and fail here.
for (auto const& [isAsync, wasFound, value, bucket] : std::array<Expected, 4>{
{{.isAsync = true, .wasFound = true, .value = 3.0, .bucket = 2},
{.isAsync = true, .wasFound = false, .value = 9.0, .bucket = 3},
{.isAsync = false, .wasFound = true, .value = 40.0, .bucket = 5},
{.isAsync = false, .wasFound = false, .value = 800.0, .bucket = 9}}})
{
auto const* point = findPoint(points, isAsync, wasFound);
ASSERT_NE(point, nullptr) << "missing series for fetch_type="
<< fetchTypeLabelValue(isAsync)
<< " found=" << fetchFoundLabelValue(wasFound);
EXPECT_EQ(point->count_, 1u);
EXPECT_EQ(opentelemetry::nostd::get<double>(point->sum_), value);
EXPECT_EQ(point->counts_[bucket], 1u);
}
}
TEST(NodeStoreReadHistogram, a_guarded_negative_latency_never_reaches_the_instrument)
{
HistogramFixture fixture;
// Negative path, end to end: the record site consults
// shouldRecordFetchLatency() before calling Record(), so a clock anomaly
// produces no sample at all. Reproduced here with the same guard.
for (double const elapsedUs : {-1.0, -1'000.0})
{
if (shouldRecordFetchLatency(static_cast<long long>(elapsedUs)))
fixture.record(elapsedUs, /*isAsync=*/false, /*wasFound=*/true);
}
// No series at all: nothing was recorded, so the instrument exported
// nothing rather than exporting a zero-count point.
EXPECT_TRUE(fixture.collect().empty());
// Cause, not just state: the same fixture DOES accept a valid sample, so
// the emptiness above is the guard working and not a broken fixture.
fixture.record(9.0, /*isAsync=*/false, /*wasFound=*/true);
auto const* point = findPoint(fixture.collect(), /*isAsync=*/false, /*wasFound=*/true);
ASSERT_NE(point, nullptr);
EXPECT_EQ(point->count_, 1u);
EXPECT_EQ(point->counts_[3], 1u);
}
#endif // XRPL_ENABLE_TELEMETRY

View File

@@ -396,11 +396,7 @@ public:
logs_->journal("JobQueue"),
*logs_,
*perfLog_))
// `*this` is passed as a ServiceRegistry only to reach the
// MetricsRegistry later, from onFetch(). It is not dereferenced here,
// and metricsRegistry_ does not exist yet at this point in the
// initializer list -- the metric macros null-check it per call.
, nodeStoreScheduler_(*this, *jobQueue_)
, nodeStoreScheduler_(*jobQueue_)
, shaMapStore_(makeSHAMapStore(*this, nodeStoreScheduler_, logs_->journal("SHAMapStore")))
, tempNodeCache_(
"NodeCache",

View File

@@ -1,29 +1,15 @@
// cspell:ignore ISTOGRAM
// The all-caps macro name XRPL_METRIC_HISTOGRAM_RECORD_LABELED trips cspell's
// compound-word splitter, which emits the subword "ISTOGRAM"; ignore it here.
#include <xrpld/app/main/NodeStoreScheduler.h>
#include <xrpld/telemetry/MetricMacros.h>
#include <xrpl/core/Job.h>
#include <xrpl/core/JobQueue.h>
#include <xrpl/core/ServiceRegistry.h>
#include <xrpl/nodestore/Scheduler.h>
#include <xrpl/nodestore/Task.h>
#include <xrpl/telemetry/NodeStoreMetricNames.h>
#include <chrono>
#include <string>
namespace xrpl {
NodeStoreScheduler::NodeStoreScheduler([[maybe_unused]] ServiceRegistry& app, JobQueue& jobQueue)
#ifdef XRPL_ENABLE_TELEMETRY
: app_(app), jobQueue_(jobQueue)
#else
: jobQueue_(jobQueue)
#endif
NodeStoreScheduler::NodeStoreScheduler(JobQueue& jobQueue) : jobQueue_(jobQueue)
{
}
@@ -47,36 +33,14 @@ NodeStoreScheduler::onFetch(node_store::FetchReport const& report)
if (jobQueue_.isStopped())
return;
auto const isAsync = report.fetchType == node_store::FetchType::Async;
// The report is in microseconds but addLoadEvents takes milliseconds, so
// cast explicitly. The load monitor only tracks whole-millisecond load,
// so the sub-millisecond detail is deliberately dropped here; the
// histogram below keeps it.
// so the sub-millisecond detail is deliberately dropped here; telemetry
// reads the microsecond value from the nodestore instead.
jobQueue_.addLoadEvents(
isAsync ? JtNsAsyncRead : JtNsSyncRead,
report.fetchType == node_store::FetchType::Async ? JtNsAsyncRead : JtNsSyncRead,
1,
std::chrono::duration_cast<std::chrono::milliseconds>(report.elapsed));
// Skip a negative elapsed time rather than hand it to the SDK, which
// rejects it and logs a warning on every single call. The clock is
// monotonic, so this needs a clock bug to happen -- but a per-fetch log
// flood would be worse than the missing sample.
if (!telemetry::shouldRecordFetchLatency(report.elapsed.count()))
return;
// Two labels, both already on the report. fetch_type because a slow async
// read only delays prefetch while a slow sync read blocks a caller;
// found because a miss can cost a read of every backend, so mixing the
// two blurs the distribution.
XRPL_METRIC_HISTOGRAM_RECORD_LABELED(
app_,
telemetry::kNodeStoreReadUs,
telemetry::kNodeStoreReadUsDesc,
report.elapsed.count(),
{{telemetry::kFetchTypeLabel, std::string(telemetry::fetchTypeLabelValue(isAsync))},
{telemetry::kFetchFoundLabel,
std::string(telemetry::fetchFoundLabelValue(report.wasFound))}});
}
void

View File

@@ -6,57 +6,13 @@
namespace xrpl {
// Forward-declared rather than included: only a reference is stored, and
// ServiceRegistry.h pulls in <boost/asio.hpp>, which every file including this
// header would then pay for. The .cpp includes the full definition.
class ServiceRegistry;
/**
* A node_store::Scheduler which uses the JobQueue.
*
* Two responsibilities, both delegating outward: it turns backend write
* requests into JobQueue jobs, and it forwards completion reports to the
* load monitor and to the OTel metrics pipeline.
*
* Collaborator diagram (ASCII):
*
* node_store::Database / Backend
* |
* | scheduleTask / onFetch / onBatchWrite
* v
* +---------------------+
* | NodeStoreScheduler |
* +---------------------+
* | |
* v v
* JobQueue ServiceRegistry
* (jobs + (-> MetricsRegistry,
* load events) read-latency histogram)
*
* @note Thread safety: onFetch() and onBatchWrite() are called from backend
* read/write threads, concurrently. Both members are references to
* objects that outlive this one, and every call they make
* (JobQueue::addLoadEvents, OTel Histogram::Record) is itself
* thread-safe, so no locking is needed here.
* @note Lifetime: constructed early, in the Application member initializer
* list, which is BEFORE the MetricsRegistry exists. It therefore stores
* the ServiceRegistry and resolves the registry per call; the metric
* macros null-check it, so fetches completing before the registry is
* created are simply not recorded.
*/
class NodeStoreScheduler : public node_store::Scheduler
{
public:
/**
* Construct a scheduler.
*
* @param app Service registry, used only to reach the
* MetricsRegistry when reporting read latency. Must
* outlive this object.
* @param jobQueue Queue that runs scheduled write tasks and receives
* load events. Must outlive this object.
*/
NodeStoreScheduler(ServiceRegistry& app, JobQueue& jobQueue);
explicit NodeStoreScheduler(JobQueue& jobQueue);
void
scheduleTask(node_store::Task& task) override;
@@ -66,21 +22,6 @@ public:
onBatchWrite(node_store::BatchWriteReport const& report) override;
private:
#ifdef XRPL_ENABLE_TELEMETRY
/**
* Service registry, resolved to a MetricsRegistry on each onFetch().
*
* Only needed when OTel is compiled in, since the record site is the
* only reader. Held under the guard for the same reason
* MetricsRegistry::app_ is: without OTel it would be an unused private
* field, which -Wall rejects.
*/
ServiceRegistry& app_;
#endif // XRPL_ENABLE_TELEMETRY
/**
* Queue used for scheduled tasks and load-event reporting.
*/
JobQueue& jobQueue_;
};

View File

@@ -52,7 +52,6 @@
#include <xrpl/server/LoadFeeTrack.h>
#include <xrpl/server/NetworkOPs.h>
#include <xrpl/telemetry/GetObjectMetricNames.h>
#include <xrpl/telemetry/NodeStoreMetricNames.h>
#include <opentelemetry/context/context.h>
#include <opentelemetry/exporters/otlp/otlp_http_metric_exporter_factory.h>
@@ -153,35 +152,6 @@ constexpr std::array kMicrosecondBoundaries{
30'000'000.0,
60'000'000.0};
/**
* Bucket boundaries for latencies that are normally sub-millisecond.
*
* 1 µs, 2 µs, 5 µs, 10 µs, 25 µs, 50 µs, 100 µs, 250 µs, 500 µs, 1 ms, 5 ms,
* 25 ms.
*
* kMicrosecondBoundaries starts at 100 µs, which is above the entire range a
* healthy nodestore read occupies, so every warm read falls in its first
* bucket and the distribution reads as flat. These edges resolve the warm
* range instead, while still reaching far enough to show a cold tail against
* it.
*
* Used by the kNodeStoreReadUs view registered in
* initExporterAndProvider().
*/
constexpr std::array kSubMillisecondBoundaries{
1.0,
2.0,
5.0,
10.0,
25.0,
50.0,
100.0,
250.0,
500.0,
1'000.0,
5'000.0,
25'000.0};
/**
* Register an explicit-bucket histogram view.
*
@@ -268,23 +238,6 @@ addRoundDurationHistogramView(metric_sdk::ViewRegistry& views, std::string const
120'000.0});
}
/**
* Register the sub-millisecond-ladder view for a duration instrument.
*
* For latencies that normally sit below one millisecond, where the
* microsecond ladder's 100 µs first edge would swallow the whole healthy
* range. See kSubMillisecondBoundaries.
*
* @param views The registry to add the view to.
* @param name Instrument name to match.
*/
void
addSubMillisecondHistogramView(metric_sdk::ViewRegistry& views, std::string const& name)
{
addHistogramView(
views, name, {kSubMillisecondBoundaries.begin(), kSubMillisecondBoundaries.end()});
}
} // namespace
#endif // XRPL_ENABLE_TELEMETRY
@@ -374,13 +327,6 @@ MetricsRegistry::initExporterAndProvider(std::string const& endpoint, std::strin
// comes from the shared constant both sites use.
addMicrosecondHistogramView(*views, kGetObjectLookupUs);
// Per-fetch nodestore read latency. Recorded at its NodeStoreScheduler.cpp
// call site, so again the name comes from the shared constant. This one
// gets the sub-millisecond ladder, not the microsecond one: a warm read
// answers in single-digit microseconds, which is below the microsecond
// ladder's first edge.
addSubMillisecondHistogramView(*views, kNodeStoreReadUs);
// Sweep malloc_trim duration. Shares the microsecond ladder rather than
// getting a bespoke one, and the ladder is what makes it readable: a trim on
// a small heap lands in the tens-of-microseconds buckets, while a trim on a
@@ -1702,7 +1648,8 @@ MetricsRegistry::registerStorageDetailGauge()
// --- Task 7.13: Storage detail gauges ---
// Reports the cumulative payload bytes handed to the NodeStore. See the
// note at the observe() call below: this is logical bytes stored, not
// on-disk file size, because no accessor for the latter exists.
// on-disk file size, because no accessor for the latter exists. The label
// value names it that way so it is not read as a filesystem measurement.
storageDetailGauge_ =
meter_->CreateInt64ObservableGauge("storage_detail", "Storage detail metrics");
storageDetailGauge_->AddCallback(
@@ -1720,11 +1667,13 @@ MetricsRegistry::registerStorageDetailGauge()
->Observe(value, {{"metric", name}});
};
// Cumulative payload bytes handed to the NodeStore -- NOT
// on-disk file size, despite the name. getStoreSize() sums
// the object payloads this process has written, so it
// excludes NuDB's keys, bucket padding and log, and it
// resets with the process while the files do not.
// Cumulative payload bytes handed to the NodeStore. This is
// not an on-disk file size: getStoreSize() sums the object
// payloads this process has written, so it excludes NuDB's
// keys, bucket padding and log, and it resets with the
// process while the files do not. The value comes from
// Database, not from any backend, so it carries no nudb_
// prefix -- it reads the same on RocksDB.
//
// This is the same call node_written_bytes makes on the
// nodestore_state gauge, so the two series are equal by
@@ -1734,7 +1683,8 @@ MetricsRegistry::registerStorageDetailGauge()
// computing one would mean stat()ing the backend's files
// from the reader thread, which needs a new Backend method
// rather than a change at this call site.
observe("nudb_bytes", static_cast<int64_t>(app.getNodeStore().getStoreSize()));
observe(
"stored_object_bytes", static_cast<int64_t>(app.getNodeStore().getStoreSize()));
}
catch (...) // NOLINT(bugprone-empty-catch)
{

View File

@@ -49,7 +49,7 @@
* +-- CountedObject counts
* +-- Load factor breakdown
* +-- NodeStore I/O gauges (totals, derived means, NuDB write queue,
* ledger-acquisition stall counters)
* ledger-acquisition stall counters)
* +-- Server info (state, uptime, peers, consensus)
* +-- Build info (version label)
* +-- Complete ledger ranges (start/end pairs)
@@ -614,6 +614,84 @@ public:
{
return meter_;
}
/**
* Sink handed to the nodestore_state gauge helpers below.
*
* Every value they publish multiplexes onto the single `nodestore_state`
* gauge through its `metric` label, so the helpers need no access to the
* OTel observer result -- just somewhere to put a name and a number.
*/
using ObserveFn = std::function<void(char const* name, std::int64_t value)>;
/**
* Observe the NodeStore I/O totals and the means derived from them.
*
* Publishes the four cumulative totals (`node_reads_total`,
* `node_writes`, `node_reads_duration_us`, `node_writes_duration_us`)
* unconditionally, plus `read_mean_us` and `write_mean_us` derived from
* them via scaledMean(). `write_mean_us` is the signal for the "a node
* with a large existing database syncs slower than a fresh one" symptom:
* back-fill is write-bound, so no read-side reading can show it. All
* three concrete store paths time themselves through
* Database::recordStoreDuration(), so the write mean is live on an
* ordinary node.
*
* Gauge rather than histogram, deliberately. A histogram would give true
* percentiles, but it costs one Record() per node object on the
* store/fetch path, and one ledger write walks thousands of SHAMap
* nodes. This reads the existing atomics once per ~10 s tick and adds
* nothing to the hot path. Consequence, stated plainly: p99 is NOT
* obtainable from this signal. A histogram added later would also need an
* explicit-bucket View registered via addMicrosecondHistogramView(),
* because the SDK's default buckets top out at 10,000.
*
* @param db NodeStore to read the counters from.
* @param observe Sink for one `metric`-labelled value.
*
* @note The totals are monotonic and never reset, so a panel wanting
* current rather than since-boot latency divides the two rates. That is
* why the counts and duration totals are exported beside the means.
* @note A mean is omitted when its count is 0, so a dashboard shows a gap
* rather than a plausible-looking 0 us.
*/
static void
observeNodeStoreTotals(node_store::Database& db, ObserveFn const& observe);
/**
* Observe the backend write-path detail, when the backend measures it.
*
* Publishes nothing for a backend whose getWriteStats() is std::nullopt,
* which is every backend except NuDB. Absent labels let a reader tell
* "not measured" from "measured, and idle"; zeros would read as a
* perfectly idle write path.
*
* @param db NodeStore whose writable backend is sampled.
* @param observe Sink for one `metric`-labelled value.
*/
static void
observeWritePathDetail(node_store::Database const& db, ObserveFn const& observe);
/**
* Observe the ledger-acquisition progress and stall counters.
*
* @param stats Process-wide acquisition counters.
* @param observe Sink for one `metric`-labelled value.
*/
static void
observeAcquireStats(AcquireStats const& stats, ObserveFn const& observe);
/**
* Observe the read queue depth and the read thread-pool counts.
*
* These four have no accessor on Database, so its JSON counters object
* is still the only way to reach them.
*
* @param db NodeStore to read the JSON counters from.
* @param observe Sink for one `metric`-labelled value.
*/
static void
observeReadQueue(node_store::Database& db, ObserveFn const& observe);
#endif
private:
@@ -961,85 +1039,11 @@ private:
void
registerNodeStoreGauge(); // Task 9.1
/**
* Sink handed to the registerNodeStoreGauge() helpers below.
*
* Every value they publish multiplexes onto the single `nodestore_state`
* gauge through its `metric` label, so the helpers need no access to the
* OTel observer result -- just somewhere to put a name and a number.
* That also makes each helper directly unit-testable by passing a
* recording sink, which the real observer result is not.
*/
using ObserveFn = std::function<void(char const* name, std::int64_t value)>;
// The four nodestore_state helpers and their ObserveFn sink are public
// (above), so a test can drive each one with a recording sink and assert
// the exact `metric` label values it publishes. They read only their
// arguments, so exposing them widens no state.
/**
* Observe the NodeStore I/O totals and the means derived from them.
*
* Publishes the four cumulative totals (`node_reads_total`,
* `node_writes`, `node_reads_duration_us`, `node_writes_duration_us`)
* unconditionally, plus `read_mean_us` and `write_mean_us` derived from
* them via scaledMean(). `write_mean_us` is the signal for the "a node
* with a large existing database syncs slower than a fresh one" symptom:
* back-fill is write-bound, so no read-side reading can show it. All
* three concrete store paths time themselves through
* Database::recordStoreDuration(), so the write mean is live on an
* ordinary node.
*
* Gauge rather than histogram, deliberately. A histogram would give true
* percentiles, but it costs one Record() per node object on the
* store/fetch path, and one ledger write walks thousands of SHAMap
* nodes. This reads the existing atomics once per ~10 s tick and adds
* nothing to the hot path. Consequence, stated plainly: p99 is NOT
* obtainable from this signal. A histogram added later would also need an
* explicit-bucket View registered via addMicrosecondHistogramView(),
* because the SDK's default buckets top out at 10,000.
*
* @param db NodeStore to read the counters from.
* @param observe Sink for one `metric`-labelled value.
*
* @note The totals are monotonic and never reset, so a panel wanting
* current rather than since-boot latency divides the two rates. That is
* why the counts and duration totals are exported beside the means.
* @note A mean is omitted when its count is 0, so a dashboard shows a gap
* rather than a plausible-looking 0 us.
*/
static void
observeNodeStoreTotals(node_store::Database& db, ObserveFn const& observe);
/**
* Observe the backend write-path detail, when the backend measures it.
*
* Publishes nothing for a backend whose getWriteStats() is std::nullopt,
* which is every backend except NuDB. Absent labels let a reader tell
* "not measured" from "measured, and idle"; zeros would read as a
* perfectly idle write path.
*
* @param db NodeStore whose writable backend is sampled.
* @param observe Sink for one `metric`-labelled value.
*/
static void
observeWritePathDetail(node_store::Database const& db, ObserveFn const& observe);
/**
* Observe the ledger-acquisition progress and stall counters.
*
* @param stats Process-wide acquisition counters.
* @param observe Sink for one `metric`-labelled value.
*/
static void
observeAcquireStats(AcquireStats const& stats, ObserveFn const& observe);
/**
* Observe the read queue depth and the read thread-pool counts.
*
* These four have no accessor on Database, so its JSON counters object
* is still the only way to reach them.
*
* @param db NodeStore to read the JSON counters from.
* @param observe Sink for one `metric`-labelled value.
*/
static void
observeReadQueue(node_store::Database& db, ObserveFn const& observe);
void
registerRotationStateGauge(); // Sync diagnostics: online_delete rotation
void