From b2bc4024b4a65aef55ec07b0bfd20eaa1db25bb0 Mon Sep 17 00:00:00 2001 From: Pratik Mankawde <3397372+pratikmankawde@users.noreply.github.com> Date: Wed, 26 Aug 2026 18:18:17 +0100 Subject: [PATCH] fix(telemetry): bound the validation tracker's pending map on insert ValidationTracker::pending_ was pruned only by evictOldPending(), which is reachable only from reconcile(), which runs only from the observable-gauge callbacks. Those callbacks need telemetry compiled in and enabled, so any node without both recorded one entry per validated ledger and freed none -- roughly 2.4 MB a day, plus a mutex acquisition per ledger. kMaxPendingEvents did not help: that check lives inside the function that never ran. This also affected ordinary builds, not just telemetry-off ones, because startAsyncGauges() returns early when [telemetry] enabled=0 and so never registers the callbacks. Bound the map where the bound can actually be guaranteed -- the insert path -- by dropping the oldest entry once it is full. A container has to be bounded by its writer, not by whoever happens to read it. Also stop doing the work when nothing will consume it: hold the tracker, its accessor and its header behind the telemetry guard, and check isEnabled() at both record sites, which the metric macros do but these direct calls did not. Adds pendingCount() so the bound is observable, and a test that records 8000 ledgers without ever reconciling and asserts the size settles at a fixed cap. --- .../libxrpl/telemetry/ValidationTracker.cpp | 36 +++++++++++++++++++ src/xrpld/app/consensus/RCLConsensus.cpp | 9 ++++- src/xrpld/app/ledger/detail/LedgerMaster.cpp | 8 ++++- src/xrpld/telemetry/MetricsRegistry.h | 20 ++++++++--- src/xrpld/telemetry/ValidationTracker.h | 31 ++++++++++++++++ .../telemetry/detail/ValidationTracker.cpp | 28 +++++++++++++++ 6 files changed, 125 insertions(+), 7 deletions(-) diff --git a/src/tests/libxrpl/telemetry/ValidationTracker.cpp b/src/tests/libxrpl/telemetry/ValidationTracker.cpp index 2542316e7a..9d176377a0 100644 --- a/src/tests/libxrpl/telemetry/ValidationTracker.cpp +++ b/src/tests/libxrpl/telemetry/ValidationTracker.cpp @@ -377,3 +377,39 @@ TEST_F(ValidationTrackerTest, GrossAgreementsCountInitialOnly) // Additive invariant: gross agree + gross miss == ledgers reconciled. EXPECT_EQ(tracker_.totalAgreementsEver() + tracker_.totalMissedEver(), 5u); } + +// --------------------------------------------------------------- +// 12. Pending map stays bounded when nothing ever reconciles +// reconcile() is the only pruning path, and it runs only from +// the observable-gauge callbacks -- which need telemetry both +// compiled in and enabled. A node with telemetry off, or with +// [telemetry] enabled=0, therefore never reconciles, so the +// record methods have to bound the map themselves. +// --------------------------------------------------------------- +TEST_F(ValidationTrackerTest, PendingStaysBoundedWithoutReconcile) +{ + constexpr std::uint64_t kFirstBatch = 4000; + constexpr std::uint64_t kSecondBatch = 8000; + + for (std::uint64_t i = 0; i < kFirstBatch; ++i) + tracker_.recordOurValidation(makeHash(i), static_cast(i)); + + auto const afterFirst = tracker_.pendingCount(); + + for (std::uint64_t i = kFirstBatch; i < kSecondBatch; ++i) + tracker_.recordOurValidation(makeHash(i), static_cast(i)); + + auto const afterSecond = tracker_.pendingCount(); + + // Bounded at all: far fewer entries retained than recorded. + EXPECT_LT(afterFirst, kFirstBatch); + + // Bounded at a fixed cap, not merely growing more slowly: doubling the + // input leaves the size unchanged. Asserted without naming the private + // constant, so the test survives a change to its value. + EXPECT_EQ(afterFirst, afterSecond); + + // Every recorded validation is still counted, so bounding the map does not + // cost us the lifetime totals the gauges report. + EXPECT_EQ(tracker_.totalValidationsSent(), kSecondBatch); +} diff --git a/src/xrpld/app/consensus/RCLConsensus.cpp b/src/xrpld/app/consensus/RCLConsensus.cpp index 9948781a7a..a5ed99dc0e 100644 --- a/src/xrpld/app/consensus/RCLConsensus.cpp +++ b/src/xrpld/app/consensus/RCLConsensus.cpp @@ -1072,9 +1072,16 @@ RCLConsensus::Adaptor::validate(RCLCxLedger const& ledger, RCLTxSet const& txns, if (auto* mr = app_.getMetricsRegistry()) { mr->incrementValidationsSent(); +#ifdef XRPL_ENABLE_TELEMETRY // Record our validation for the agreement tracker so it can // compare against network-validated ledgers. - mr->getValidationTracker().recordOurValidation(ledger.id(), ledger.seq()); + // + // Only when enabled: recording takes the tracker's lock and inserts an + // entry, and nothing reconciles or drains those entries unless the + // observable gauges are running. + if (mr->isEnabled()) + mr->getValidationTracker().recordOurValidation(ledger.id(), ledger.seq()); +#endif } } diff --git a/src/xrpld/app/ledger/detail/LedgerMaster.cpp b/src/xrpld/app/ledger/detail/LedgerMaster.cpp index 63c00ef8e0..cf8d65230b 100644 --- a/src/xrpld/app/ledger/detail/LedgerMaster.cpp +++ b/src/xrpld/app/ledger/detail/LedgerMaster.cpp @@ -296,10 +296,16 @@ LedgerMaster::setValidLedger(std::shared_ptr const& l) (void)maxLedgerDifference_; validLedgerSeq_ = l->header().seq; +#ifdef XRPL_ENABLE_TELEMETRY // Record the network-validated ledger for the agreement tracker so it // can compare against our own validations. - if (auto* mr = app_.getMetricsRegistry()) + // + // Only when enabled: recording takes the tracker's lock and inserts an + // entry, and nothing reconciles or drains those entries unless the + // observable gauges are running. + if (auto* mr = app_.getMetricsRegistry(); mr && mr->isEnabled()) mr->getValidationTracker().recordNetworkValidation(l->header().hash, l->header().seq); +#endif app_.getOPs().updateLocalTx(*l); app_.getSHAMapStore().onLedgerClosed(getValidatedLedger()); diff --git a/src/xrpld/telemetry/MetricsRegistry.h b/src/xrpld/telemetry/MetricsRegistry.h index 3ca31af991..0837bd13ce 100644 --- a/src/xrpld/telemetry/MetricsRegistry.h +++ b/src/xrpld/telemetry/MetricsRegistry.h @@ -138,7 +138,11 @@ * instrumentation site. */ +#ifdef XRPL_ENABLE_TELEMETRY +// The tracker is held and exposed only in this configuration, where the gauge +// callbacks that drain it exist. #include +#endif #include @@ -655,10 +659,15 @@ public: void incrementTxqDropped(std::string_view reason); +#ifdef XRPL_ENABLE_TELEMETRY /** * Access the validation agreement tracker. * Used by consensus and ledger hooks to record our validations and * network validations so the tracker can compute agreement percentages. + * + * Guarded, along with the tracker itself, because only the observable-gauge + * callbacks read it and those exist only in this configuration. Recording + * into it is not free: each call takes its lock and inserts an entry. * @return Reference to the internal ValidationTracker instance. */ ValidationTracker& @@ -667,7 +676,6 @@ public: return validationTracker_; } -#ifdef XRPL_ENABLE_TELEMETRY /** * Access the shared OTel Meter for call-site instrument creation. * Used by the XRPL_METRIC_* macros (MetricMacros.h) so new synchronous @@ -742,15 +750,17 @@ private: */ bool const enabled_; +#ifdef XRPL_ENABLE_TELEMETRY /** * Tracks validation agreement between this node and the network. - * Lives outside the XRPL_ENABLE_TELEMETRY guard because it is - * always safe to record events; the gauge callback simply won't - * fire when telemetry is disabled. + * + * Guarded because reconcile() -- which resolves and then prunes recorded + * events -- runs only from the observable-gauge callbacks. Recording + * without it accumulates one entry per validated ledger, so the tracker + * exists only where something drains it. */ ValidationTracker validationTracker_; -#ifdef XRPL_ENABLE_TELEMETRY /** * Reference to Application services for gauge callbacks. * Only needed when OTel is compiled in, since observable gauge diff --git a/src/xrpld/telemetry/ValidationTracker.h b/src/xrpld/telemetry/ValidationTracker.h index ac80f5cdf5..61f711602f 100644 --- a/src/xrpld/telemetry/ValidationTracker.h +++ b/src/xrpld/telemetry/ValidationTracker.h @@ -253,6 +253,17 @@ public: uint64_t totalValidationsChecked() const; + /** + * Number of ledgers currently held awaiting reconciliation. + * + * Never exceeds kMaxPendingEvents: the record methods enforce that bound + * as they insert, so the map stays bounded whether or not anything ever + * reconciles or reads it. + * @return Size of the pending map. + */ + [[nodiscard]] std::size_t + pendingCount() const; + /** @} */ private: @@ -397,6 +408,26 @@ private: void evictOldPending(TimePoint now); + /** + * Hold pending_ at kMaxPendingEvents by dropping its oldest entry. + * + * Called on the insert path, because that is the only place the bound can + * be guaranteed. reconcile() also prunes, but it runs only while the gauge + * callbacks are registered, which needs telemetry both compiled in and + * enabled -- so a node with telemetry off, or with [telemetry] enabled=0, + * would otherwise grow this map by one entry per validated ledger forever. + * + * Drops the oldest entry rather than the least useful one: the map is + * unordered, so this is a linear scan, but it runs at most once per + * recorded validation and only once the map is already full. + * + * @param justRecorded Hash inserted by the caller, kept even if the scan + * finds it oldest (equal timestamps make that possible). + * @note Caller must hold mutex_. + */ + void + boundPending(uint256 const& justRecorded); + /** * Scan a window deque and flip the first non-agreed entry matching * the given ledger hash to agreed. diff --git a/src/xrpld/telemetry/detail/ValidationTracker.cpp b/src/xrpld/telemetry/detail/ValidationTracker.cpp index c7f9c599bc..89e74c8f56 100644 --- a/src/xrpld/telemetry/detail/ValidationTracker.cpp +++ b/src/xrpld/telemetry/detail/ValidationTracker.cpp @@ -31,6 +31,7 @@ ValidationTracker::recordOurValidation(uint256 const& ledgerHash, LedgerIndex se } evt.weValidated = true; totalValidationsSent_.fetch_add(1, std::memory_order_relaxed); + boundPending(ledgerHash); } void @@ -46,6 +47,33 @@ ValidationTracker::recordNetworkValidation(uint256 const& ledgerHash, LedgerInde } evt.networkValidated = true; totalValidationsChecked_.fetch_add(1, std::memory_order_relaxed); + boundPending(ledgerHash); +} + +void +ValidationTracker::boundPending(uint256 const& justRecorded) +{ + if (pending_.size() <= kMaxPendingEvents) + return; + + auto oldest = pending_.end(); + for (auto it = pending_.begin(); it != pending_.end(); ++it) + { + if (it->first == justRecorded) + continue; + if (oldest == pending_.end() || it->second.recordTime < oldest->second.recordTime) + oldest = it; + } + + if (oldest != pending_.end()) + pending_.erase(oldest); +} + +std::size_t +ValidationTracker::pendingCount() const +{ + std::scoped_lock const lock(mutex_); + return pending_.size(); } void