From d655cf24e99c4142055e67c39ef87f953aa5a736 Mon Sep 17 00:00:00 2001 From: Nicholas Dudfield Date: Fri, 18 Sep 2026 11:57:07 +0700 Subject: [PATCH] feat(overlay): weld PeerImp I/O onto Transport Inbound and outbound peers wrap the ssl_stream in SslTransport. SimOverlay can now inject an in-process transport without a PeerImp template. --- src/xrpld/overlay/detail/ConnectAttempt.cpp | 2 +- src/xrpld/overlay/detail/OverlayImpl.cpp | 2 +- src/xrpld/overlay/detail/PeerImp.cpp | 116 +++++++++----------- src/xrpld/overlay/detail/PeerImp.h | 19 ++-- 4 files changed, 64 insertions(+), 75 deletions(-) diff --git a/src/xrpld/overlay/detail/ConnectAttempt.cpp b/src/xrpld/overlay/detail/ConnectAttempt.cpp index 4bb76369c4..af1d6cfc25 100644 --- a/src/xrpld/overlay/detail/ConnectAttempt.cpp +++ b/src/xrpld/overlay/detail/ConnectAttempt.cpp @@ -382,7 +382,7 @@ ConnectAttempt::processResponse() auto const peer = std::make_shared( app_, - std::move(stream_ptr_), + std::make_unique(std::move(stream_ptr_)), read_buf_.data(), std::move(slot_), std::move(response_), diff --git a/src/xrpld/overlay/detail/OverlayImpl.cpp b/src/xrpld/overlay/detail/OverlayImpl.cpp index 3768616dc6..9a03e1ff14 100644 --- a/src/xrpld/overlay/detail/OverlayImpl.cpp +++ b/src/xrpld/overlay/detail/OverlayImpl.cpp @@ -291,7 +291,7 @@ OverlayImpl::onHandoff( publicKey, *negotiatedVersion, consumer, - std::move(stream_ptr), + std::make_unique(std::move(stream_ptr)), *this); { // As we are not on the strand, run() must be called diff --git a/src/xrpld/overlay/detail/PeerImp.cpp b/src/xrpld/overlay/detail/PeerImp.cpp index 88301721ef..5dfdd548a7 100644 --- a/src/xrpld/overlay/detail/PeerImp.cpp +++ b/src/xrpld/overlay/detail/PeerImp.cpp @@ -73,7 +73,7 @@ PeerImp::PeerImp( PublicKey const& publicKey, ProtocolVersion protocol, Resource::Consumer consumer, - std::unique_ptr&& stream_ptr, + std::unique_ptr&& transport, OverlayImpl& overlay) : Child(overlay) , app_(app) @@ -82,11 +82,9 @@ PeerImp::PeerImp( , p_sink_(app_.journal("Protocol"), makePrefix(id)) , journal_(sink_) , p_journal_(p_sink_) - , stream_ptr_(std::move(stream_ptr)) - , socket_(stream_ptr_->next_layer().socket()) - , stream_(*stream_ptr_) - , strand_(makePeerStrand(app.config(), socket_.get_executor())) - , timer_(waitable_timer{socket_.get_executor()}) + , transport_(std::move(transport)) + , strand_(makePeerStrand(app.config(), transport_->get_executor())) + , timer_(waitable_timer{transport_->get_executor()}) , vtimer_(app.makePeerTimer()) , remote_address_(slot->remote_endpoint()) , overlay_(overlay) @@ -219,7 +217,7 @@ PeerImp::stop() { if (!strand_.running_in_this_thread()) return post(strand_, std::bind(&PeerImp::stop, shared_from_this())); - if (socket_.is_open()) + if (transport_->is_open()) { // The rationale for using different severity levels is that // outbound connections are under our control and may be logged @@ -281,17 +279,15 @@ PeerImp::send(std::shared_ptr const& m) if (sendq_size != 0) return; - boost::asio::async_write( - stream_, - boost::asio::buffer( - send_queue_.front()->getBuffer(compressionEnabled_)), - bind_executor( - strand_, - std::bind( - &PeerImp::onWriteMessage, - shared_from_this(), - std::placeholders::_1, - std::placeholders::_2))); + transport_->async_write( + toConstBuffers(boost::asio::buffer( + send_queue_.front()->getBuffer(compressionEnabled_))), + strand_, + std::bind( + &PeerImp::onWriteMessage, + shared_from_this(), + std::placeholders::_1, + std::placeholders::_2)); } void @@ -572,12 +568,12 @@ PeerImp::close() XRPL_ASSERT( strand_.running_in_this_thread(), "ripple::PeerImp::close : strand in this thread"); - if (socket_.is_open()) + if (transport_->is_open()) { detaching_ = true; // DEPRECATED error_code ec; timer_.cancel(ec); - socket_.close(ec); + transport_->close(); overlay_.incPeerDisconnect(); if (inbound_) { @@ -600,7 +596,7 @@ PeerImp::fail(std::string const& reason) (void(Peer::*)(std::string const&)) & PeerImp::fail, shared_from_this(), reason)); - if (journal_.active(beast::severities::kWarning) && socket_.is_open()) + if (journal_.active(beast::severities::kWarning) && transport_->is_open()) { std::string const n = name(); JLOG(journal_.warn()) << (n.empty() ? remote_address_.to_string() : n) @@ -615,7 +611,7 @@ PeerImp::fail(std::string const& name, error_code ec) XRPL_ASSERT( strand_.running_in_this_thread(), "ripple::PeerImp::fail : strand in this thread"); - if (socket_.is_open()) + if (transport_->is_open()) { JLOG(journal_.warn()) << name << " from " << toBase58(TokenType::NodePublic, publicKey_) @@ -631,7 +627,8 @@ PeerImp::gracefulClose() strand_.running_in_this_thread(), "ripple::PeerImp::gracefulClose : strand in this thread"); XRPL_ASSERT( - socket_.is_open(), "ripple::PeerImp::gracefulClose : socket is open"); + transport_->is_open(), + "ripple::PeerImp::gracefulClose : socket is open"); XRPL_ASSERT( !gracefulClose_, "ripple::PeerImp::gracefulClose : socket is not closing"); @@ -639,10 +636,10 @@ PeerImp::gracefulClose() if (send_queue_.size() > 0) return; setTimer(); - stream_.async_shutdown(bind_executor( + transport_->async_shutdown( strand_, std::bind( - &PeerImp::onShutdown, shared_from_this(), std::placeholders::_1))); + &PeerImp::onShutdown, shared_from_this(), std::placeholders::_1)); } void @@ -696,7 +693,7 @@ PeerImp::makePrefix(id_t id) void PeerImp::onTimer(error_code const& ec) { - if (!socket_.is_open()) + if (!transport_->is_open()) return; if (ec == boost::asio::error::operation_aborted) @@ -779,7 +776,7 @@ PeerImp::doAccept() JLOG(journal_.debug()) << "doAccept: " << remote_address_; - auto const sharedValue = makeSharedValue(*stream_ptr_, journal_); + auto const sharedValue = transport_->makeSharedValue(journal_); // This shouldn't fail since we already computed // the shared value successfully in OverlayImpl @@ -818,15 +815,12 @@ PeerImp::doAccept() app_); // Write the whole buffer and only start protocol when that's done. - boost::asio::async_write( - stream_, - write_buffer->data(), - boost::asio::transfer_all(), - bind_executor( - strand_, - [this, write_buffer, self = shared_from_this()]( - error_code ec, std::size_t bytes_transferred) { - if (!socket_.is_open()) + transport_->async_write( + toConstBuffers(write_buffer->data()), + strand_, + [this, write_buffer, self = shared_from_this()]( + error_code ec, std::size_t bytes_transferred) { + if (!transport_->is_open()) return; if (ec == boost::asio::error::operation_aborted) return; @@ -835,7 +829,7 @@ PeerImp::doAccept() if (write_buffer->size() == bytes_transferred) return doProtocolStart(); return fail("Failed to write header"); - })); + }); } std::string @@ -896,7 +890,7 @@ PeerImp::doProtocolStart() void PeerImp::onReadMessage(error_code ec, std::size_t bytes_transferred) { - if (!socket_.is_open()) + if (!transport_->is_open()) return; if (ec == boost::asio::error::operation_aborted) return; @@ -936,7 +930,7 @@ PeerImp::onReadMessage(error_code ec, std::size_t bytes_transferred) if (ec) return fail("onReadMessage", ec); - if (!socket_.is_open()) + if (!transport_->is_open()) return; if (gracefulClose_) return; @@ -946,21 +940,21 @@ PeerImp::onReadMessage(error_code ec, std::size_t bytes_transferred) } // Timeout on writes only - stream_.async_read_some( - read_buffer_.prepare(std::max(Tuning::readBufferBytes, hint)), - bind_executor( - strand_, - std::bind( - &PeerImp::onReadMessage, - shared_from_this(), - std::placeholders::_1, - std::placeholders::_2))); + transport_->async_read_some( + toMutableBuffers( + read_buffer_.prepare(std::max(Tuning::readBufferBytes, hint))), + strand_, + std::bind( + &PeerImp::onReadMessage, + shared_from_this(), + std::placeholders::_1, + std::placeholders::_2)); } void PeerImp::onWriteMessage(error_code ec, std::size_t bytes_transferred) { - if (!socket_.is_open()) + if (!transport_->is_open()) return; if (ec == boost::asio::error::operation_aborted) return; @@ -983,27 +977,25 @@ PeerImp::onWriteMessage(error_code ec, std::size_t bytes_transferred) if (!send_queue_.empty()) { // Timeout on writes only - return boost::asio::async_write( - stream_, - boost::asio::buffer( - send_queue_.front()->getBuffer(compressionEnabled_)), - bind_executor( - strand_, - std::bind( - &PeerImp::onWriteMessage, - shared_from_this(), - std::placeholders::_1, - std::placeholders::_2))); + return transport_->async_write( + toConstBuffers(boost::asio::buffer( + send_queue_.front()->getBuffer(compressionEnabled_))), + strand_, + std::bind( + &PeerImp::onWriteMessage, + shared_from_this(), + std::placeholders::_1, + std::placeholders::_2)); } if (gracefulClose_) { - return stream_.async_shutdown(bind_executor( + return transport_->async_shutdown( strand_, std::bind( &PeerImp::onShutdown, shared_from_this(), - std::placeholders::_1))); + std::placeholders::_1)); } } diff --git a/src/xrpld/overlay/detail/PeerImp.h b/src/xrpld/overlay/detail/PeerImp.h index 4797a92d18..64a4707647 100644 --- a/src/xrpld/overlay/detail/PeerImp.h +++ b/src/xrpld/overlay/detail/PeerImp.h @@ -27,6 +27,7 @@ #include #include #include +#include #include #include #include @@ -76,9 +77,7 @@ private: beast::WrappedSink p_sink_; beast::Journal const journal_; beast::Journal const p_journal_; - std::unique_ptr stream_ptr_; - socket_type& socket_; - stream_type& stream_; + std::unique_ptr transport_; boost::asio::strand strand_; waitable_timer timer_; std::unique_ptr vtimer_; @@ -248,7 +247,7 @@ public: PublicKey const& publicKey, ProtocolVersion protocol, Resource::Consumer consumer, - std::unique_ptr&& stream_ptr, + std::unique_ptr&& transport, OverlayImpl& overlay); /** Create outgoing, handshaked peer. */ @@ -256,7 +255,7 @@ public: template PeerImp( Application& app, - std::unique_ptr&& stream_ptr, + std::unique_ptr&& transport, Buffers const& buffers, std::shared_ptr&& slot, http_response_type&& response, @@ -666,7 +665,7 @@ private: template PeerImp::PeerImp( Application& app, - std::unique_ptr&& stream_ptr, + std::unique_ptr&& transport, Buffers const& buffers, std::shared_ptr&& slot, http_response_type&& response, @@ -682,11 +681,9 @@ PeerImp::PeerImp( , p_sink_(app_.journal("Protocol"), makePrefix(id)) , journal_(sink_) , p_journal_(p_sink_) - , stream_ptr_(std::move(stream_ptr)) - , socket_(stream_ptr_->next_layer().socket()) - , stream_(*stream_ptr_) - , strand_(makePeerStrand(app.config(), socket_.get_executor())) - , timer_(waitable_timer{socket_.get_executor()}) + , transport_(std::move(transport)) + , strand_(makePeerStrand(app.config(), transport_->get_executor())) + , timer_(waitable_timer{transport_->get_executor()}) , vtimer_(app.makePeerTimer()) , remote_address_(slot->remote_endpoint()) , overlay_(overlay)