File indexing completed on 2026-08-17 09:00:18
0001
0002
0003
0004
0005
0006
0007 #ifndef BOOST_REDIS_REDIS_STREAM_HPP
0008 #define BOOST_REDIS_REDIS_STREAM_HPP
0009
0010 #include <boost/redis/config.hpp>
0011 #include <boost/redis/detail/connect_fsm.hpp>
0012 #include <boost/redis/error.hpp>
0013 #include <boost/redis/logger.hpp>
0014
0015 #include <boost/asio/basic_waitable_timer.hpp>
0016 #include <boost/asio/cancel_after.hpp>
0017 #include <boost/asio/compose.hpp>
0018 #include <boost/asio/connect.hpp>
0019 #include <boost/asio/coroutine.hpp>
0020 #include <boost/asio/ip/basic_resolver.hpp>
0021 #include <boost/asio/ip/tcp.hpp>
0022 #include <boost/asio/local/stream_protocol.hpp>
0023 #include <boost/asio/ssl/context.hpp>
0024 #include <boost/asio/ssl/stream.hpp>
0025 #include <boost/asio/ssl/stream_base.hpp>
0026 #include <boost/asio/steady_timer.hpp>
0027 #include <boost/system/error_code.hpp>
0028
0029 #include <utility>
0030
0031 namespace boost {
0032 namespace redis {
0033 namespace detail {
0034
0035 template <class Executor>
0036 class redis_stream {
0037 asio::ssl::context ssl_ctx_;
0038 asio::ip::basic_resolver<asio::ip::tcp, Executor> resolv_;
0039 asio::ssl::stream<asio::basic_stream_socket<asio::ip::tcp, Executor>> stream_;
0040 #ifdef BOOST_ASIO_HAS_LOCAL_SOCKETS
0041 asio::basic_stream_socket<asio::local::stream_protocol, Executor> unix_socket_;
0042 #endif
0043 typename asio::steady_timer::template rebind_executor<Executor>::other timer_;
0044 redis_stream_state st_;
0045
0046 void reset_stream() { stream_ = {resolv_.get_executor(), ssl_ctx_}; }
0047
0048 struct connect_op {
0049 redis_stream& obj_;
0050 connect_fsm fsm_;
0051
0052 template <class Self>
0053 void execute_action(Self& self, connect_action act)
0054 {
0055 auto& obj = this->obj_;
0056 const auto& cfg = fsm_.get_config();
0057
0058 switch (act.type) {
0059 case connect_action_type::unix_socket_close:
0060 #ifdef BOOST_ASIO_HAS_LOCAL_SOCKETS
0061 {
0062 system::error_code ec;
0063 obj.unix_socket_.close(ec);
0064 (*this)(self, ec);
0065 }
0066 #else
0067 BOOST_ASSERT(false);
0068 #endif
0069 return;
0070 case connect_action_type::unix_socket_connect:
0071 #ifdef BOOST_ASIO_HAS_LOCAL_SOCKETS
0072 obj.unix_socket_.async_connect(
0073 cfg.unix_socket,
0074 asio::cancel_after(obj.timer_, cfg.connect_timeout, std::move(self)));
0075 #else
0076 BOOST_ASSERT(false);
0077 #endif
0078 return;
0079
0080 case connect_action_type::tcp_resolve:
0081 obj.resolv_.async_resolve(
0082 cfg.addr.host,
0083 cfg.addr.port,
0084 asio::cancel_after(obj.timer_, cfg.resolve_timeout, std::move(self)));
0085 return;
0086 case connect_action_type::ssl_stream_reset:
0087 obj.reset_stream();
0088
0089 (*this)(self);
0090 return;
0091 case connect_action_type::ssl_handshake:
0092 obj.stream_.async_handshake(
0093 asio::ssl::stream_base::client,
0094 asio::cancel_after(obj.timer_, cfg.ssl_handshake_timeout, std::move(self)));
0095 return;
0096 case connect_action_type::done: self.complete(act.ec); break;
0097
0098 case connect_action_type::tcp_connect:
0099 default: BOOST_ASSERT(false);
0100 }
0101 }
0102
0103
0104 template <class Self>
0105 void operator()(
0106 Self& self,
0107 system::error_code ec,
0108 const asio::ip::tcp::endpoint& selected_endpoint)
0109 {
0110 auto act = fsm_.resume(
0111 ec,
0112 selected_endpoint,
0113 obj_.st_,
0114 self.get_cancellation_state().cancelled());
0115 execute_action(self, act);
0116 }
0117
0118
0119 template <class Self>
0120 void operator()(
0121 Self& self,
0122 system::error_code ec,
0123 asio::ip::tcp::resolver::results_type endpoints)
0124 {
0125 auto act = fsm_.resume(ec, endpoints, obj_.st_, self.get_cancellation_state().cancelled());
0126 if (act.type == connect_action_type::tcp_connect) {
0127 auto& obj = this->obj_;
0128 asio::async_connect(
0129 obj.stream_.next_layer(),
0130 std::move(endpoints),
0131 asio::cancel_after(obj.timer_, fsm_.get_config().connect_timeout, std::move(self)));
0132 } else {
0133 execute_action(self, act);
0134 }
0135 }
0136
0137 template <class Self>
0138 void operator()(Self& self, system::error_code ec = {})
0139 {
0140 auto act = fsm_.resume(ec, obj_.st_, self.get_cancellation_state().cancelled());
0141 execute_action(self, act);
0142 }
0143 };
0144
0145 public:
0146 explicit redis_stream(Executor ex, asio::ssl::context&& ssl_ctx)
0147 : ssl_ctx_{std::move(ssl_ctx)}
0148 , resolv_{ex}
0149 , stream_{ex, ssl_ctx_}
0150 #ifdef BOOST_ASIO_HAS_LOCAL_SOCKETS
0151 , unix_socket_{ex}
0152 #endif
0153 , timer_{std::move(ex)}
0154 { }
0155
0156
0157 using executor_type = Executor;
0158 executor_type get_executor() noexcept { return resolv_.get_executor(); }
0159
0160
0161 const auto& get_ssl_context() const noexcept { return ssl_ctx_; }
0162 bool is_open() const
0163 {
0164 #ifdef BOOST_ASIO_HAS_LOCAL_SOCKETS
0165 if (st_.type == transport_type::unix_socket)
0166 return unix_socket_.is_open();
0167 #endif
0168 return stream_.next_layer().is_open();
0169 }
0170 auto& next_layer() { return stream_; }
0171 const auto& next_layer() const { return stream_; }
0172
0173
0174 template <class CompletionToken>
0175 auto async_connect(const config& cfg, buffered_logger& l, CompletionToken&& token)
0176 {
0177 return asio::async_compose<CompletionToken, void(system::error_code)>(
0178 connect_op{*this, connect_fsm(cfg, l)},
0179 token);
0180 }
0181
0182
0183 template <class ConstBufferSequence, class CompletionToken>
0184 void async_write_some(const ConstBufferSequence& buffers, CompletionToken&& token)
0185 {
0186 switch (st_.type) {
0187 case transport_type::tcp:
0188 {
0189 stream_.next_layer().async_write_some(buffers, std::forward<CompletionToken>(token));
0190 break;
0191 }
0192 case transport_type::tcp_tls:
0193 {
0194 stream_.async_write_some(buffers, std::forward<CompletionToken>(token));
0195 break;
0196 }
0197 #ifdef BOOST_ASIO_HAS_LOCAL_SOCKETS
0198 case transport_type::unix_socket:
0199 {
0200 unix_socket_.async_write_some(buffers, std::forward<CompletionToken>(token));
0201 break;
0202 }
0203 #endif
0204 default: BOOST_ASSERT(false);
0205 }
0206 }
0207
0208 template <class MutableBufferSequence, class CompletionToken>
0209 void async_read_some(const MutableBufferSequence& buffers, CompletionToken&& token)
0210 {
0211 switch (st_.type) {
0212 case transport_type::tcp:
0213 {
0214 return stream_.next_layer().async_read_some(
0215 buffers,
0216 std::forward<CompletionToken>(token));
0217 break;
0218 }
0219 case transport_type::tcp_tls:
0220 {
0221 return stream_.async_read_some(buffers, std::forward<CompletionToken>(token));
0222 break;
0223 }
0224 #ifdef BOOST_ASIO_HAS_LOCAL_SOCKETS
0225 case transport_type::unix_socket:
0226 {
0227 unix_socket_.async_read_some(buffers, std::forward<CompletionToken>(token));
0228 break;
0229 }
0230 #endif
0231 default: BOOST_ASSERT(false);
0232 }
0233 }
0234
0235
0236
0237
0238
0239 void cancel_resolve() { resolv_.cancel(); }
0240 };
0241
0242 }
0243 }
0244 }
0245
0246 #endif