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