Clio  develop
The XRP Ledger API server.
Loading...
Searching...
No Matches
WsConnection.hpp
1#pragma once
2
3#include "util/OverloadSet.hpp"
4#include "util/Taggable.hpp"
5#include "util/build/Build.hpp"
6#include "web/ng/Connection.hpp"
7#include "web/ng/Error.hpp"
8#include "web/ng/Request.hpp"
9#include "web/ng/Response.hpp"
10#include "web/ng/impl/Concepts.hpp"
11#include "web/ng/impl/SendingQueue.hpp"
12
13#include <boost/asio/buffer.hpp>
14#include <boost/asio/ip/tcp.hpp>
15#include <boost/asio/spawn.hpp>
16#include <boost/asio/ssl/context.hpp>
17#include <boost/asio/ssl/stream.hpp>
18#include <boost/beast/core/buffers_to_string.hpp>
19#include <boost/beast/core/flat_buffer.hpp>
20#include <boost/beast/core/role.hpp>
21#include <boost/beast/core/tcp_stream.hpp>
22#include <boost/beast/http/field.hpp>
23#include <boost/beast/http/message.hpp>
24#include <boost/beast/http/string_body.hpp>
25#include <boost/beast/ssl.hpp>
26#include <boost/beast/websocket/rfc6455.hpp>
27#include <boost/beast/websocket/stream.hpp>
28#include <boost/beast/websocket/stream_base.hpp>
29
30#include <chrono>
31#include <cstddef>
32#include <memory>
33#include <string>
34#include <utility>
35#include <variant>
36
37namespace web::ng::impl {
38
40public:
42
43 virtual std::expected<void, Error>
44 sendShared(std::shared_ptr<std::string> message, boost::asio::yield_context yield) = 0;
45};
46
47template <typename StreamType>
48class WsConnection : public WsConnectionBase {
49 boost::beast::websocket::stream<StreamType> stream_;
50 boost::beast::http::request<boost::beast::http::string_body> initialRequest_;
51
52 using MessageType = std::variant<Response, std::shared_ptr<std::string>>;
53 SendingQueue<MessageType> sendingQueue_;
54
55 bool closed_{false};
56
57public:
58 WsConnection(
59 StreamType&& stream,
60 std::string ip,
61 boost::beast::flat_buffer buffer,
62 boost::beast::http::request<boost::beast::http::string_body> initialRequest,
63 util::TagDecoratorFactory const& tagDecoratorFactory,
64 size_t maxSendingQueueSize
65 )
66 : WsConnectionBase(std::move(ip), std::move(buffer), tagDecoratorFactory)
67 , stream_(std::move(stream))
68 , initialRequest_(std::move(initialRequest))
69 , sendingQueue_{
70 [this](MessageType const& message, auto&& yield) {
71 boost::asio::const_buffer const buffer = std::visit(
73 [](Response const& r) -> boost::asio::const_buffer {
74 return r.asWsResponse();
75 },
76 [](std::shared_ptr<std::string> const& m) -> boost::asio::const_buffer {
77 return boost::asio::buffer(*m);
78 }
79 },
80 message
81 );
82 stream_.async_write(buffer, yield);
83 },
84 maxSendingQueueSize
85 }
86 {
87 setupWsStream();
88 }
89
90 ~WsConnection() override = default;
91 WsConnection(WsConnection&&) = delete;
92 WsConnection&
93 operator=(WsConnection&&) = delete;
94 WsConnection(WsConnection const&) = delete;
95 WsConnection&
96 operator=(WsConnection const&) = delete;
97
98 std::expected<void, Error>
99 performHandshake(boost::asio::yield_context yield)
100 {
101 Error error;
102 stream_.async_accept(initialRequest_, yield[error]);
103 if (error)
104 return std::unexpected{error};
105 return {};
106 }
107
108 [[nodiscard]] bool
109 wasUpgraded() const override
110 {
111 return true;
112 }
113
114 std::expected<void, Error>
115 sendShared(std::shared_ptr<std::string> message, boost::asio::yield_context yield) override
116 {
117 return sendingQueue_.send(std::move(message), yield);
118 }
119
120 void
121 setTimeout(std::chrono::steady_clock::duration newTimeout) override
122 {
123 boost::beast::websocket::stream_base::timeout wsTimeout =
124 boost::beast::websocket::stream_base::timeout::suggested(
125 boost::beast::role_type::server
126 );
127 wsTimeout.idle_timeout = newTimeout;
128 wsTimeout.handshake_timeout = newTimeout;
129 stream_.set_option(wsTimeout);
130 }
131
132 std::expected<void, Error>
133 send(Response response, boost::asio::yield_context yield) override
134 {
135 return sendingQueue_.send(std::move(response), yield);
136 }
137
138 std::expected<Request, Error>
139 receive(boost::asio::yield_context yield) override
140 {
141 Error error;
142 stream_.async_read(buffer_, yield[error]);
143 if (error)
144 return std::unexpected{error};
145
146 auto request = boost::beast::buffers_to_string(buffer_.data());
147 buffer_.consume(buffer_.size());
148
149 return Request{std::move(request), initialRequest_};
150 }
151
152 void
153 close(boost::asio::yield_context yield) override
154 {
155 if (closed_)
156 return;
157
158 // This should be set before the async_close(). Otherwise there is a possibility to have
159 // multiple coroutines waiting on async_close(), but only one will be woken up after the
160 // actual close happened, others will hang.
161 closed_ = true;
162
163 boost::system::error_code error; // unused
164 stream_.async_close(boost::beast::websocket::close_code::normal, yield[error]);
165 }
166
167private:
168 void
169 setupWsStream()
170 {
171 // Disable the timeout. The websocket::stream uses its own timeout settings.
172 boost::beast::get_lowest_layer(stream_).expires_never();
174 stream_.set_option(
175 boost::beast::websocket::stream_base::decorator(
176 [](boost::beast::websocket::response_type& res) {
177 res.set(
178 boost::beast::http::field::server, util::build::getClioFullVersionString()
179 );
180 }
181 )
182 );
183 }
184};
185
186using PlainWsConnection = WsConnection<boost::beast::tcp_stream>;
187using SslWsConnection = WsConnection<boost::asio::ssl::stream<boost::beast::tcp_stream>>;
188
189template <typename StreamType>
190std::expected<std::unique_ptr<WsConnection<StreamType>>, Error>
191makeWsConnection(
192 StreamType&& stream,
193 std::string ip,
194 boost::beast::flat_buffer buffer,
195 boost::beast::http::request<boost::beast::http::string_body> request,
196 util::TagDecoratorFactory const& tagDecoratorFactory,
197 size_t maxSendingQueueSize,
198 boost::asio::yield_context yield
199)
200{
201 auto connection = std::make_unique<WsConnection<StreamType>>(
202 std::forward<StreamType>(stream),
203 std::move(ip),
204 std::move(buffer),
205 std::move(request),
206 tagDecoratorFactory,
207 maxSendingQueueSize
208 );
209 auto const expectedSuccess = connection->performHandshake(yield);
210 if (not expectedSuccess.has_value())
211 return std::unexpected{expectedSuccess.error()};
212 return connection;
213}
214
215} // namespace web::ng::impl
A factory for TagDecorator instantiation.
Definition Taggable.hpp:165
std::string const & ip() const
Get the ip of the client.
Definition Connection.cpp:21
static constexpr std::chrono::steady_clock::duration kDefaultTimeout
The default timeout for send, receive, and close operations.
Definition Connection.hpp:124
Connection(std::string ip, boost::beast::flat_buffer buffer, util::TagDecoratorFactory const &tagDecoratorFactory)
Construct a new Connection object.
Definition Connection.cpp:32
Represents an HTTP or WebSocket request.
Definition Request.hpp:18
Represents an HTTP or Websocket response.
Definition Response.hpp:21
Definition SendingQueue.hpp:19
Definition WsConnection.hpp:39
Connection(std::string ip, boost::beast::flat_buffer buffer, util::TagDecoratorFactory const &tagDecoratorFactory)
Construct a new Connection object.
Definition Connection.cpp:32
std::expected< Request, Error > receive(boost::asio::yield_context yield) override
Receive a request from the client.
Definition WsConnection.hpp:139
void close(boost::asio::yield_context yield) override
Gracefully close the connection.
Definition WsConnection.hpp:153
std::expected< void, Error > send(Response response, boost::asio::yield_context yield) override
Send a response to the client.
Definition WsConnection.hpp:133
void setTimeout(std::chrono::steady_clock::duration newTimeout) override
Get the timeout for send, receive, and close operations. For WebSocket connections,...
Definition WsConnection.hpp:121
bool wasUpgraded() const override
Whether the connection was upgraded. Upgraded connections are websocket connections.
Definition WsConnection.hpp:109
Overload set for lambdas.
Definition OverloadSet.hpp:11