From 9f92a938e51d107136a24fc498a3a662c91f08a9 Mon Sep 17 00:00:00 2001 From: Pratik Mankawde <3397372+pratikmankawde@users.noreply.github.com> Date: Tue, 28 Jul 2026 16:08:56 +0100 Subject: [PATCH] test(telemetry): assert the nodestore_state labels and pin writer depth The 22 metric label values on the nodestore_state gauge were asserted nowhere, and the depthSum assertions could not tell the real fetch_add(depth) from a fetch_add(1) that would silently zero every derived queueing time. Make the four observe* helpers and their ObserveFn sink public rather than private, so the seam the header already claimed is actually reachable. All four are static and read only their arguments, so this exposes no object state; a friend declaration would have granted access to every private member instead. The assertions live in the existing Beast nodestore suite because src/test compiles into xrpld, which contains MetricsRegistry.cpp, while xrpl_tests deliberately does not when telemetry is enabled. Each helper now has its exact emitted label set asserted, so a typo in any literal fails instead of silently producing a disjoint series, and each derived mean is asserted ABSENT on a fresh store -- a refactor to value_or(0) would draw a believable flat zero on a latency axis and otherwise pass everything. Replace the concurrent write-stats test with one that forces genuine overlap through a latch. NuDB holds one global mutex for the whole insert and doInsert reads the depth before entering it, so a blocked thread records a depth of at least 2; asserting depthSum strictly exceeds insertCount therefore cannot be satisfied by a constant 1. The old bounds admitted that bug at their floor. Also: cover the std::nullopt branch on the two backends that exist in every build, bound the store-duration accumulator by the wall clock, drop four assertions that cannot fail, and correct two comments that claimed coverage the tests do not have -- the duplicate-key test is not the throwing path, because nudb reports key_exists without throwing, and no test drives NodeStoreScheduler::onFetch. --- .../scripts/levelization/results/ordering.txt | 2 + src/test/nodestore/DatabaseConfig_test.cpp | 513 ++++++++++++++++-- src/tests/libxrpl/nodestore/Backend.cpp | 50 ++ src/tests/libxrpl/nodestore/Database.cpp | 15 +- src/tests/libxrpl/nodestore/NuDBFactory.cpp | 206 +++++-- .../libxrpl/telemetry/MetricsRegistry.cpp | 3 +- .../telemetry/NodeStoreMetricNames.cpp | 28 +- src/xrpld/telemetry/MetricsRegistry.h | 112 ++-- 8 files changed, 795 insertions(+), 134 deletions(-) diff --git a/.github/scripts/levelization/results/ordering.txt b/.github/scripts/levelization/results/ordering.txt index 494fbc25e3..d9c4f24ecf 100644 --- a/.github/scripts/levelization/results/ordering.txt +++ b/.github/scripts/levelization/results/ordering.txt @@ -138,7 +138,9 @@ test.nodestore > test.jtx test.nodestore > test.unit_test test.nodestore > xrpl.basics test.nodestore > xrpl.config +test.nodestore > xrpld.app test.nodestore > xrpld.core +test.nodestore > xrpld.telemetry test.nodestore > xrpl.nodestore test.nodestore > xrpl.rdb test.overlay > test.jtx diff --git a/src/test/nodestore/DatabaseConfig_test.cpp b/src/test/nodestore/DatabaseConfig_test.cpp index 5c9b2cf0d6..b8392daaa8 100644 --- a/src/test/nodestore/DatabaseConfig_test.cpp +++ b/src/test/nodestore/DatabaseConfig_test.cpp @@ -4,6 +4,13 @@ #include #include +#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 +#include +#endif #include #include @@ -21,16 +28,21 @@ #include #include #include +#include #include #include #include +#include +#include #include #include -#include #include +#include +#include #include #include +#include namespace xrpl::node_store { @@ -624,12 +636,6 @@ public: std::is_same_v().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::max() > std::numeric_limits::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(kNumStored), 4321); BEAST_EXPECT(stored.size() == kNumStored); + + auto const writesBegan = std::chrono::steady_clock::now(); storeBatch(*db, stored); + auto const wallClockUs = + static_cast(std::chrono::duration_cast( + 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 + 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 writableBackend = + Manager::instance().makeBackend(writableParams, megabytes(4), scheduler, journal_); + std::shared_ptr archiveBackend = + Manager::instance().makeBackend(archiveParams, megabytes(4), scheduler, journal_); + if (!writableBackend || !archiveBackend) + return nullptr; + writableBackend->open(); + archiveBackend->open(); + + return std::make_unique( + 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 writableBackend = - Manager::instance().makeBackend(writableParams, megabytes(4), scheduler, journal_); - std::shared_ptr 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> 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 + names() const + { + std::vector 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 + 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 + 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 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 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(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(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 total, + std::optional count) { + return (total && count && *count != 0) ? std::optional{*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 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 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 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(kNumStored), 97531); + storeBatch(*db, stored); + + MetricSink busy; + telemetry::MetricsRegistry::observeWritePathDetail(*db, busy.fn()); + std::vector 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 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 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(); } }; diff --git a/src/tests/libxrpl/nodestore/Backend.cpp b/src/tests/libxrpl/nodestore/Backend.cpp index cc1d515621..c1f5b95e84 100644 --- a/src/tests/libxrpl/nodestore/Backend.cpp +++ b/src/tests/libxrpl/nodestore/Backend.cpp @@ -204,4 +204,54 @@ INSTANTIATE_TEST_SUITE_P( ::testing::ValuesIn(backendTypes()), [](::testing::TestParamInfo 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 diff --git a/src/tests/libxrpl/nodestore/Database.cpp b/src/tests/libxrpl/nodestore/Database.cpp index 298c906a23..2f89442815 100644 --- a/src/tests/libxrpl/nodestore/Database.cpp +++ b/src/tests/libxrpl/nodestore/Database.cpp @@ -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); } diff --git a/src/tests/libxrpl/nodestore/NuDBFactory.cpp b/src/tests/libxrpl/nodestore/NuDBFactory.cpp index 1153aa3749..2658cd3fdf 100644 --- a/src/tests/libxrpl/nodestore/NuDBFactory.cpp +++ b/src/tests/libxrpl/nodestore/NuDBFactory.cpp @@ -4,6 +4,7 @@ #include #include #include +#include #include #include @@ -14,7 +15,9 @@ #include #include #include +#include #include +#include #include #include #include @@ -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(kOverlapThreads)); + + std::vector 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(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 threads; - threads.reserve(kThreads); - for (auto t = 0uz; t < kThreads; ++t) + std::uint64_t completedRounds = 0; + std::optional 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(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(); } diff --git a/src/tests/libxrpl/telemetry/MetricsRegistry.cpp b/src/tests/libxrpl/telemetry/MetricsRegistry.cpp index 490d622ab1..f21ddfe831 100644 --- a/src/tests/libxrpl/telemetry/MetricsRegistry.cpp +++ b/src/tests/libxrpl/telemetry/MetricsRegistry.cpp @@ -411,7 +411,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); } diff --git a/src/tests/libxrpl/telemetry/NodeStoreMetricNames.cpp b/src/tests/libxrpl/telemetry/NodeStoreMetricNames.cpp index 0c558ef252..22f1bebeb6 100644 --- a/src/tests/libxrpl/telemetry/NodeStoreMetricNames.cpp +++ b/src/tests/libxrpl/telemetry/NodeStoreMetricNames.cpp @@ -25,9 +25,26 @@ * 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. + * bucket ladder the value is filed into. + * + * KNOWN GAP, stated plainly. Two separate steps remain unverified by any test: + * + * - That NodeStoreScheduler::onFetch() calls the histogram macro at all, with + * report.elapsed.count() and those two labels. `NodeStoreScheduler` appears + * in exactly one test file, src/test/app/SHAMapStore_test.cpp, and only as a + * construction site that asserts nothing about onFetch. The nodestore + * suites use DummyScheduler or a test-local CapturingScheduler, neither of + * which is NodeStoreScheduler, so they do NOT cover this call path. + * - That the recorded value has sub-millisecond resolution end to end. That + * half IS covered, in src/tests/libxrpl/nodestore/Database.cpp's + * sub_millisecond_fetch_latency_is_reported, which drives a real store + * through a capturing Scheduler and asserts the reported microsecond total + * equals the nodestore's own accumulator exactly. + * + * Closing the first gap needs an integration test under src/test/ with a live + * jtx::Env and an in-memory metric reader attached to the running + * MetricsRegistry; that is larger than this file's scope and is deliberately + * deferred rather than claimed. */ #include @@ -84,7 +101,6 @@ TEST(NodeStoreMetricNames, instrument_name_is_the_exact_shared_literal) // 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 @@ -160,10 +176,6 @@ TEST(NodeStoreMetricNames, latency_guard_admits_zero_and_refuses_negatives) 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)); } // --------------------------------------------------------------------------- diff --git a/src/xrpld/telemetry/MetricsRegistry.h b/src/xrpld/telemetry/MetricsRegistry.h index d6dab25e12..6c10cdafc6 100644 --- a/src/xrpld/telemetry/MetricsRegistry.h +++ b/src/xrpld/telemetry/MetricsRegistry.h @@ -50,7 +50,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) @@ -612,6 +612,59 @@ 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; + + /** + * Observe the NodeStore I/O totals and the means derived from them. + * + * @param db NodeStore to read the counters from. + * @param observe Sink for one `metric`-labelled value. + */ + 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: @@ -894,60 +947,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; + // 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. - * - * @param db NodeStore to read the counters from. - * @param observe Sink for one `metric`-labelled value. - */ - 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 registerServerInfoGauge(); // Task 9.7a void