From 645f16e7fe495100a35bcb119d2bcbd36a09d7c0 Mon Sep 17 00:00:00 2001 From: Nicholas Dudfield Date: Fri, 18 Sep 2026 11:47:01 +0700 Subject: [PATCH] feat(overlay): virtual peer timer, inline strands, precise clock PeerImp arms TimeoutCounterTimer when makePeerTimer is installed. inlineStrands builds the existing strand type over an inline executor. Tracking, ping RTT and uptime use getPreciseStopwatch. --- src/xrpld/overlay/detail/PeerImp.cpp | 31 +++++++--- src/xrpld/overlay/detail/PeerImp.h | 20 ++++-- src/xrpld/overlay/detail/PeerStrand.h | 87 +++++++++++++++++++++++++++ 3 files changed, 125 insertions(+), 13 deletions(-) create mode 100644 src/xrpld/overlay/detail/PeerStrand.h diff --git a/src/xrpld/overlay/detail/PeerImp.cpp b/src/xrpld/overlay/detail/PeerImp.cpp index b7319862d1..88301721ef 100644 --- a/src/xrpld/overlay/detail/PeerImp.cpp +++ b/src/xrpld/overlay/detail/PeerImp.cpp @@ -40,6 +40,7 @@ #include #include +#include #include #include @@ -84,17 +85,18 @@ PeerImp::PeerImp( , stream_ptr_(std::move(stream_ptr)) , socket_(stream_ptr_->next_layer().socket()) , stream_(*stream_ptr_) - , strand_(socket_.get_executor()) + , strand_(makePeerStrand(app.config(), socket_.get_executor())) , timer_(waitable_timer{socket_.get_executor()}) + , vtimer_(app.makePeerTimer()) , remote_address_(slot->remote_endpoint()) , overlay_(overlay) , inbound_(true) , protocol_(protocol) , tracking_(Tracking::unknown) - , trackingTime_(clock_type::now()) + , trackingTime_(steadyNow()) , publicKey_(publicKey) - , lastPingTime_(clock_type::now()) - , creationTime_(clock_type::now()) + , lastPingTime_(steadyNow()) + , creationTime_(steadyNow()) , squelch_(app_.journal("Squelch")) , usage_(consumer) , fee_{Resource::feeTrivialPeer, ""} @@ -646,6 +648,14 @@ PeerImp::gracefulClose() void PeerImp::setTimer() { + if (vtimer_) + { + vtimer_->expiresAfter(peerTimerInterval, [self = shared_from_this()]() { + boost::asio::dispatch( + self->strand_, [self]() { self->onTimer({}); }); + }); + return; + } error_code ec; timer_.expires_from_now(peerTimerInterval, ec); @@ -664,6 +674,11 @@ PeerImp::setTimer() void PeerImp::cancelTimer() { + if (vtimer_) + { + vtimer_->cancel(); + return; + } error_code ec; timer_.cancel(ec); } @@ -706,7 +721,7 @@ PeerImp::onTimer(error_code const& ec) { std::lock_guard sl(recentLock_); - duration = clock_type::now() - trackingTime_; + duration = steadyNow() - trackingTime_; } if ((t == Tracking::diverged && @@ -727,7 +742,7 @@ PeerImp::onTimer(error_code const& ec) return; } - lastPingTime_ = clock_type::now(); + lastPingTime_ = steadyNow(); lastPingSeq_ = rand_int(app_.getPrng()); protocol::TMPing message; @@ -1090,7 +1105,7 @@ PeerImp::onMessage(std::shared_ptr const& m) // Update latency estimate auto const rtt = std::chrono::round( - clock_type::now() - lastPingTime_); + steadyNow() - lastPingTime_); std::lock_guard sl(recentLock_); @@ -1963,7 +1978,7 @@ PeerImp::checkTracking(std::uint32_t seq1, std::uint32_t seq2) std::lock_guard sl(recentLock_); tracking_ = Tracking::diverged; - trackingTime_ = clock_type::now(); + trackingTime_ = steadyNow(); } } diff --git a/src/xrpld/overlay/detail/PeerImp.h b/src/xrpld/overlay/detail/PeerImp.h index 23916db217..4797a92d18 100644 --- a/src/xrpld/overlay/detail/PeerImp.h +++ b/src/xrpld/overlay/detail/PeerImp.h @@ -23,7 +23,9 @@ #include #include #include +#include #include +#include #include #include #include @@ -79,6 +81,7 @@ private: stream_type& stream_; boost::asio::strand strand_; waitable_timer timer_; + std::unique_ptr vtimer_; // Updated at each stage of the connection process to reflect // the current conditions as closely as possible. @@ -366,7 +369,7 @@ public: clock_type::duration uptime() const { - return clock_type::now() - creationTime_; + return steadyNow() - creationTime_; } Json::Value @@ -458,6 +461,12 @@ private: void cancelTimer(); + clock_type::time_point + steadyNow() const + { + return app_.getPreciseStopwatch().now(); + } + static std::string makePrefix(id_t id); @@ -676,17 +685,18 @@ PeerImp::PeerImp( , stream_ptr_(std::move(stream_ptr)) , socket_(stream_ptr_->next_layer().socket()) , stream_(*stream_ptr_) - , strand_(socket_.get_executor()) + , strand_(makePeerStrand(app.config(), socket_.get_executor())) , timer_(waitable_timer{socket_.get_executor()}) + , vtimer_(app.makePeerTimer()) , remote_address_(slot->remote_endpoint()) , overlay_(overlay) , inbound_(false) , protocol_(protocol) , tracking_(Tracking::unknown) - , trackingTime_(clock_type::now()) + , trackingTime_(steadyNow()) , publicKey_(publicKey) - , lastPingTime_(clock_type::now()) - , creationTime_(clock_type::now()) + , lastPingTime_(steadyNow()) + , creationTime_(steadyNow()) , squelch_(app_.journal("Squelch")) , usage_(usage) , fee_{Resource::feeTrivialPeer} diff --git a/src/xrpld/overlay/detail/PeerStrand.h b/src/xrpld/overlay/detail/PeerStrand.h new file mode 100644 index 0000000000..612d18e5e7 --- /dev/null +++ b/src/xrpld/overlay/detail/PeerStrand.h @@ -0,0 +1,87 @@ +#ifndef RIPPLE_OVERLAY_PEERSTRAND_H_INCLUDED +#define RIPPLE_OVERLAY_PEERSTRAND_H_INCLUDED + +#include +#include +#include +#include +#include + +namespace ripple { + +class InlineExecutor +{ + boost::asio::execution_context* context_; + +public: + explicit InlineExecutor(boost::asio::execution_context& ctx) noexcept + : context_(&ctx) + { + } + + [[nodiscard]] boost::asio::execution_context& + context() const noexcept + { + return *context_; + } + + void + on_work_started() const noexcept + { + } + + void + on_work_finished() const noexcept + { + } + + template + void + dispatch(F&& f, A const&) const + { + std::forward(f)(); + } + + template + void + post(F&& f, A const&) const + { + std::forward(f)(); + } + + template + void + defer(F&& f, A const&) const + { + std::forward(f)(); + } + + friend bool + operator==(InlineExecutor const& a, InlineExecutor const& b) noexcept + { + return a.context_ == b.context_; + } + + friend bool + operator!=(InlineExecutor const& a, InlineExecutor const& b) noexcept + { + return a.context_ != b.context_; + } +}; + +using PeerStrand = boost::asio::strand; + +[[nodiscard]] inline PeerStrand +makePeerStrand( + Config const& config, + boost::asio::any_io_executor const& transportExecutor) +{ + if (config.inlineStrands) + return boost::asio::make_strand(boost::asio::executor( + InlineExecutor{transportExecutor.context()})); + return boost::asio::make_strand(transportExecutor); +} + +} // namespace ripple + +#endif