1#include <test/jtx/WSClient.h>
3#include <xrpld/core/Config.h>
5#include <xrpl/basics/Mutex.hpp>
6#include <xrpl/basics/contract.h>
7#include <xrpl/config/BasicConfig.h>
8#include <xrpl/config/Constants.h>
9#include <xrpl/json/json_reader.h>
10#include <xrpl/json/json_value.h>
11#include <xrpl/json/to_string.h>
12#include <xrpl/protocol/jss.h>
13#include <xrpl/server/Port.h>
15#include <boost/asio/bind_executor.hpp>
16#include <boost/asio/buffer.hpp>
17#include <boost/asio/executor_work_guard.hpp>
18#include <boost/asio/io_context.hpp>
19#include <boost/asio/ip/address_v4.hpp>
20#include <boost/asio/ip/address_v6.hpp>
21#include <boost/asio/ip/tcp.hpp>
22#include <boost/asio/post.hpp>
23#include <boost/asio/strand.hpp>
24#include <boost/beast/core/multi_buffer.hpp>
25#include <boost/beast/websocket/error.hpp>
26#include <boost/beast/websocket/rfc6455.hpp>
27#include <boost/beast/websocket/stream.hpp>
28#include <boost/beast/websocket/stream_base.hpp>
29#include <boost/system/detail/error_code.hpp>
30#include <boost/system/system_error.hpp>
63 static boost::asio::ip::tcp::endpoint
69 auto const ps = v2 ?
"ws2" :
"ws";
78 using namespace boost::asio::ip;
79 if (pp.
ip && pp.
ip->is_unspecified())
81 *pp.
ip = pp.
ip->is_v6() ? address{address_v6::loopback()}
82 : address{address_v4::loopback()};
94 template <
class ConstBuffers>
98 using boost::asio::buffer;
99 using boost::asio::buffer_size;
102 buffer_copy(buffer(&s[0], s.
size()), b);
108 boost::asio::strand<boost::asio::io_context::executor_type>
strand_;
111 boost::beast::websocket::stream<boost::asio::ip::tcp::socket&>
ws_;
112 boost::beast::multi_buffer
rb_;
134 boost::asio::bind_executor(
strand_, [
this] {
145 catch (boost::system::system_error
const&)
152 work_ = std::nullopt;
174 boost::beast::websocket::stream_base::decorator(
175 [&](boost::beast::websocket::request_type& req) {
176 for (
auto const& h : headers)
177 req.set(h.first, h.second);
185 catch (std::exception&)
200 using boost::asio::buffer;
201 using namespace std::chrono_literals;
209 jp[jss::method] = cmd;
210 jp[jss::jsonrpc] =
"2.0";
211 jp[jss::ripplerpc] =
"2.0";
216 jp[jss::command] = cmd;
224 ws_.write_some(
true, buffer(s), ec);
230 findMsg(5s, [&](
json::Value const& jval) {
return jval[jss::type] == jss::response; });
234 jv->removeMember(jss::type);
235 if ((*jv).isMember(jss::status) && (*jv)[jss::status] == jss::error)
238 ret[jss::result] = *jv;
239 if ((*jv).isMember(jss::error))
240 ret[jss::error] = (*jv)[jss::error];
241 ret[jss::status] = jss::error;
244 if ((*jv).isMember(jss::status) && (*jv).isMember(jss::result))
245 (*jv)[jss::result][jss::status] = (*jv)[jss::status];
257 if (!
cv_.wait_for(lock, timeout, [&] { return !msgs_.empty(); }))
259 m = std::move(
msgs_.back());
262 return std::move(m->jv);
272 if (!
cv_.wait_for(lock, timeout, [&] {
273 for (auto it = msgs_.begin(); it != msgs_.end(); ++it)
288 return std::move(m->jv);
291 [[nodiscard]]
unsigned
306 boost::asio::bind_executor(
312 boost::beast::websocket::close_code::normal,
326 boost::asio::bind_executor(
329 boost::system::error_code ec;
341 if (ec == boost::beast::websocket::error::closed)
Unserialize a JSON document into a Value.
bool parse(std::string const &document, Value &root)
Read a Value from a JSON document.
Holds unparsed configuration information.
bool exists(std::string const &name) const
Returns true if a section with the given name exists.
Section & section(std::string const &name)
Returns the section with the given name.
std::vector< std::string > const & values() const
Returns all the values in the section.
std::list< std::shared_ptr< Msg > > msgs_
xrpl::Mutex< bool > readEnded_
void disconnect() override
Close the client's connection to the server.
void onReadMsg(error_code const &ec)
boost::beast::multi_buffer rb_
boost::asio::io_context ios_
unsigned version() const override
Get RPC 1.0 or RPC 2.0.
std::optional< boost::asio::executor_work_guard< boost::asio::io_context::executor_type > > work_
boost::asio::ip::tcp::socket stream_
std::optional< json::Value > getMsg(std::chrono::milliseconds const &timeout) override
Retrieve a message.
boost::beast::websocket::stream< boost::asio::ip::tcp::socket & > ws_
json::Value invoke(std::string const &cmd, json::Value const ¶ms) override
Submit a command synchronously.
std::optional< json::Value > findMsg(std::chrono::milliseconds const &timeout, std::function< bool(json::Value const &)> pred) override
Retrieve a message that meets the predicate criteria.
WSClientImpl(Config const &cfg, bool v2, unsigned rpcVersion, std::unordered_map< std::string, std::string > const &headers={})
std::condition_variable cv_
static constexpr auto kDisconnectTimeout
boost::system::error_code error_code
static std::string bufferString(ConstBuffers const &b)
static boost::asio::ip::tcp::endpoint getEndpoint(BasicConfig const &cfg, bool v2)
boost::asio::strand< boost::asio::io_context::executor_type > strand_
std::condition_variable readEndCv_
std::unique_ptr< WSClient > makeWSClient(Config const &cfg, bool v2, unsigned rpcVersion, std::unordered_map< std::string, std::string > const &headers)
Returns a client operating through WebSockets/S.
void parsePort(ParsedPort &port, Section const §ion, std::ostream &log)
std::string to_string(BaseUInt< Bits, Tag > const &a)
XRPL_NO_SANITIZE_ADDRESS void rethrow()
Rethrow the exception currently being handled.
XRPL_NO_SANITIZE_ADDRESS void Throw(Args &&... args)
std::set< std::string, boost::beast::iless > protocol
std::optional< boost::asio::ip::address > ip
std::optional< std::uint16_t > port
static constexpr auto kServer