#include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include namespace xrpl::test { class WSClientImpl : public WSClient { using error_code = boost::system::error_code; struct Msg { json::Value jv; explicit Msg(json::Value&& jv) : jv(std::move(jv)) { } }; static boost::asio::ip::tcp::endpoint getEndpoint(BasicConfig const& cfg, bool v2) { auto& log = std::cerr; ParsedPort common; parsePort(common, cfg[Sections::kServer], log); auto const ps = v2 ? "ws2" : "ws"; for (auto const& name : cfg.section(Sections::kServer).values()) { if (!cfg.exists(name)) continue; ParsedPort pp; parsePort(pp, cfg[name], log); if (!pp.protocol.contains(ps)) continue; using namespace boost::asio::ip; if (pp.ip && pp.ip->is_unspecified()) { *pp.ip = pp.ip->is_v6() ? address{address_v6::loopback()} : address{address_v4::loopback()}; } if (!pp.port) Throw("Use fixConfigPorts with auto ports"); return {*pp.ip, *pp.port}; // NOLINT(bugprone-unchecked-optional-access) } Throw("Missing WebSocket port"); return {}; // Silence compiler control paths return value warning } template static std::string bufferString(ConstBuffers const& b) { using boost::asio::buffer; using boost::asio::buffer_size; std::string s; s.resize(buffer_size(b)); buffer_copy(buffer(&s[0], s.size()), b); return s; } boost::asio::io_context ios_; std::optional> work_; boost::asio::strand strand_; std::thread thread_; boost::asio::ip::tcp::socket stream_; boost::beast::websocket::stream ws_; boost::beast::multi_buffer rb_; bool peerClosed_ = false; // disconnect() waits on this until the read loop ends (for any reason: // the server acknowledged our close, or a timeout force-closed the socket). static constexpr auto kDisconnectTimeout = std::chrono::seconds{1}; xrpl::Mutex readEnded_; std::condition_variable readEndCv_; // synchronize message queue std::mutex m_; std::condition_variable cv_; std::list> msgs_; unsigned rpcVersion_; void cleanup() { boost::asio::post( ios_, // boost::asio::bind_executor(strand_, [this] { if (!peerClosed_) { ws_.async_close( {}, // boost::asio::bind_executor(strand_, [&](error_code) { try { stream_.cancel(); } // NOLINTNEXTLINE(bugprone-empty-catch) catch (boost::system::system_error const&) { // ignored } })); } })); work_ = std::nullopt; thread_.join(); } public: WSClientImpl( Config const& cfg, bool v2, unsigned rpcVersion, std::unordered_map const& headers = {}) : work_(std::in_place, boost::asio::make_work_guard(ios_)) , strand_(boost::asio::make_strand(ios_)) , thread_([&] { ios_.run(); }) , stream_(ios_) , ws_(stream_) , rpcVersion_(rpcVersion) { try { auto const ep = getEndpoint(cfg, v2); stream_.connect(ep); ws_.set_option( boost::beast::websocket::stream_base::decorator( [&](boost::beast::websocket::request_type& req) { for (auto const& h : headers) req.set(h.first, h.second); })); ws_.handshake(ep.address().to_string() + ":" + std::to_string(ep.port()), "/"); ws_.async_read( rb_, boost::asio::bind_executor(strand_, [this](error_code const& ec, std::size_t) { onReadMsg(ec); })); } catch (std::exception&) { cleanup(); rethrow(); } } ~WSClientImpl() override { cleanup(); } json::Value invoke(std::string const& cmd, json::Value const& params) override { using boost::asio::buffer; using namespace std::chrono_literals; { json::Value jp; if (params) jp = params; if (rpcVersion_ == 2) { jp[jss::method] = cmd; jp[jss::jsonrpc] = "2.0"; jp[jss::ripplerpc] = "2.0"; jp[jss::id] = 5; } else { jp[jss::command] = cmd; } auto const s = to_string(jp); // Use the error_code overload to avoid an unhandled exception // when the server closes the WebSocket connection (e.g. after // booting a client that exceeded resource thresholds). error_code ec; ws_.write_some(true, buffer(s), ec); if (ec) return {}; } auto jv = findMsg(5s, [&](json::Value const& jval) { return jval[jss::type] == jss::response; }); if (jv) { // Normalize JSON output jv->removeMember(jss::type); if ((*jv).isMember(jss::status) && (*jv)[jss::status] == jss::error) { json::Value ret; ret[jss::result] = *jv; if ((*jv).isMember(jss::error)) ret[jss::error] = (*jv)[jss::error]; ret[jss::status] = jss::error; return ret; } if ((*jv).isMember(jss::status) && (*jv).isMember(jss::result)) (*jv)[jss::result][jss::status] = (*jv)[jss::status]; return *jv; } return {}; } std::optional getMsg(std::chrono::milliseconds const& timeout) override { std::shared_ptr m; { std::unique_lock lock(m_); if (!cv_.wait_for(lock, timeout, [&] { return !msgs_.empty(); })) return std::nullopt; m = std::move(msgs_.back()); msgs_.pop_back(); } return std::move(m->jv); } std::optional findMsg(std::chrono::milliseconds const& timeout, std::function pred) override { std::shared_ptr m; { std::unique_lock lock(m_); if (!cv_.wait_for(lock, timeout, [&] { for (auto it = msgs_.begin(); it != msgs_.end(); ++it) { if (pred((*it)->jv)) { m = std::move(*it); msgs_.erase(it); return true; } } return false; })) { return std::nullopt; } } return std::move(m->jv); } [[nodiscard]] unsigned version() const override { return rpcVersion_; } void disconnect() override { // Perform a graceful WebSocket closing handshake and block until the // read loop ends, so the server observes a clean close (not a RST) and // has finished tearing the connection down by the time we return. // If the server already closed, the wait below returns immediately. boost::asio::post( ios_, boost::asio::bind_executor( strand_, // [this] { if (!peerClosed_) { ws_.async_close( boost::beast::websocket::close_code::normal, boost::asio::bind_executor(strand_, [](error_code) {})); } })); auto lock = readEnded_.lock(); readEndCv_.wait_for(lock, kDisconnectTimeout, [&lock] { return *lock; }); // On timeout (server gone or not replying) force the socket closed so // the outstanding read ends and the worker thread can later be joined. if (!*lock) { boost::asio::post( ios_, boost::asio::bind_executor( strand_, // [this] { boost::system::error_code ec; stream_.close(ec); })); } } private: void onReadMsg(error_code const& ec) { if (ec) { if (ec == boost::beast::websocket::error::closed) peerClosed_ = true; *readEnded_.lock() = true; readEndCv_.notify_all(); return; } json::Value jv; json::Reader jr; jr.parse(bufferString(rb_.data()), jv); rb_.consume(rb_.size()); auto m = std::make_shared(std::move(jv)); { std::scoped_lock const lock(m_); msgs_.push_front(m); cv_.notify_all(); } ws_.async_read( rb_, boost::asio::bind_executor(strand_, [this](error_code const& ec, std::size_t) { onReadMsg(ec); })); } }; std::unique_ptr makeWSClient( Config const& cfg, bool v2, unsigned rpcVersion, std::unordered_map const& headers) { return std::make_unique(cfg, v2, rpcVersion, headers); } } // namespace xrpl::test