Back to home page

EIC code displayed by LXR

 
 

    


File indexing completed on 2026-08-17 09:00:18

0001 /* Copyright (c) 2018-2025 Marcelo Zimbres Silva (mzimbres@gmail.com),
0002  * Ruben Perez Hidalgo (rubenperez038 at gmail dot com)
0003  *
0004  * Distributed under the Boost Software License, Version 1.0. (See
0005  * accompanying file LICENSE.txt)
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_;  // prevent use-after-move errors
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);  // This is a sync action
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                // this action does not require yielding. Execute the next action immediately
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             // Connect should use the specialized handler, where resolver results are available
0098             case connect_action_type::tcp_connect:
0099             default:                               BOOST_ASSERT(false);
0100          }
0101       }
0102 
0103       // This overload will be used for connects
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       // This overload will be used for resolves
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_;  // prevent use-after-free errors
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    // Executor. Required to satisfy the AsyncStream concept
0157    using executor_type = Executor;
0158    executor_type get_executor() noexcept { return resolv_.get_executor(); }
0159 
0160    // Accessors
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    // I/O
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    // These functions should only be used with callbacks (e.g. within async_compose function bodies)
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    // Cancels resolve operations. Resolve operations don't support per-operation
0236    // cancellation, but resolvers have a cancel() function. Resolve operations are
0237    // in general blocking and run in a separate thread. cancel() has effect only
0238    // if the operation hasn't started yet. Still, trying is better than nothing
0239    void cancel_resolve() { resolv_.cancel(); }
0240 };
0241 
0242 }  // namespace detail
0243 }  // namespace redis
0244 }  // namespace boost
0245 
0246 #endif