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  *
0003  * Distributed under the Boost Software License, Version 1.0. (See
0004  * accompanying file LICENSE.txt)
0005  */
0006 
0007 #ifndef BOOST_REDIS_MULTIPLEXER_HPP
0008 #define BOOST_REDIS_MULTIPLEXER_HPP
0009 
0010 #include <boost/redis/adapter/adapt.hpp>
0011 #include <boost/redis/adapter/any_adapter.hpp>
0012 #include <boost/redis/config.hpp>
0013 #include <boost/redis/detail/read_buffer.hpp>
0014 #include <boost/redis/resp3/node.hpp>
0015 #include <boost/redis/resp3/parser.hpp>
0016 #include <boost/redis/resp3/type.hpp>
0017 #include <boost/redis/usage.hpp>
0018 
0019 #include <boost/system/error_code.hpp>
0020 
0021 #include <cstddef>
0022 #include <deque>
0023 #include <functional>
0024 #include <memory>
0025 #include <string_view>
0026 #include <utility>
0027 
0028 namespace boost::redis {
0029 
0030 class request;
0031 
0032 namespace detail {
0033 
0034 // Return type of the multiplexer::consume_next function
0035 enum class consume_result
0036 {
0037    needs_more,    // consume_next didn't have enough data
0038    got_response,  // got a response to a regular command, vs. a push
0039    got_push,      // got a response to a push
0040 };
0041 
0042 class multiplexer {
0043 public:
0044    struct elem {
0045    public:
0046       explicit elem(request const& req, any_adapter adapter);
0047 
0048       void set_done_callback(std::function<void()> f) noexcept { done_ = std::move(f); };
0049 
0050       auto notify_done() noexcept -> void
0051       {
0052          status_ = status::done;
0053          done_();
0054       }
0055 
0056       auto notify_error(system::error_code ec) noexcept -> void;
0057 
0058       [[nodiscard]]
0059       auto is_waiting() const noexcept
0060       {
0061          return status_ == status::waiting;
0062       }
0063 
0064       [[nodiscard]]
0065       auto is_written() const noexcept
0066       {
0067          return status_ == status::written;
0068       }
0069 
0070       [[nodiscard]]
0071       auto is_staged() const noexcept
0072       {
0073          return status_ == status::staged;
0074       }
0075 
0076       [[nodiscard]]
0077       bool is_done() const noexcept
0078       {
0079          return status_ == status::done;
0080       }
0081 
0082       void mark_written() noexcept { status_ = status::written; }
0083 
0084       void mark_staged() noexcept { status_ = status::staged; }
0085 
0086       void mark_waiting() noexcept { status_ = status::waiting; }
0087 
0088       auto get_error() const -> system::error_code const& { return ec_; }
0089 
0090       auto get_request() const -> request const& { return *req_; }
0091 
0092       auto get_read_size() const -> std::size_t { return read_size_; }
0093 
0094       auto get_remaining_responses() const -> std::size_t { return remaining_responses_; }
0095 
0096       auto commit_response(std::size_t read_size) -> void;
0097 
0098       auto get_adapter() -> any_adapter& { return adapter_; }
0099 
0100       // Marks the element as an abandoned request. An abandoned request
0101       // won't cause problems when its response arrives, but that response will be ignored.
0102       void mark_abandoned();
0103 
0104       [[nodiscard]]
0105       bool is_abandoned() const
0106       {
0107          return !req_;
0108       }
0109 
0110    private:
0111       enum class status
0112       {
0113          waiting,  // the request hasn't been written yet
0114          staged,   // we've issued the write for this request, but it hasn't finished yet
0115          written,  // the request has been written successfully
0116          done,     // the request has completed and the done callback has been invoked
0117       };
0118 
0119       request const* req_;
0120       any_adapter adapter_;
0121       std::function<void()> done_;
0122 
0123       // Contains the number of commands that haven't been read yet.
0124       std::size_t remaining_responses_;
0125       status status_;
0126 
0127       system::error_code ec_;
0128       std::size_t read_size_;
0129    };
0130 
0131    multiplexer();
0132 
0133    // To be called before a write operation. Coalesces all available requests
0134    // into a single buffer. Returns the number of coalesced requests.
0135    // Must be called before cancel_on_conn_lost() because it might change
0136    // request status.
0137    [[nodiscard]]
0138    auto prepare_write() -> std::size_t;
0139 
0140    // To be called after a write operation.
0141    // Returns true once all the bytes in the buffer generated by prepare_write
0142    // have been written.
0143    // Must be called before cancel_on_conn_lost() because it might change
0144    // request status.
0145    auto commit_write(std::size_t bytes_written) -> bool;
0146 
0147    // To be called after a successful read operation.
0148    // Must be called before cancel_on_conn_lost() because it might change
0149    // request status.
0150    [[nodiscard]]
0151    auto consume(system::error_code& ec) -> std::pair<consume_result, std::size_t>;
0152 
0153    auto add(std::shared_ptr<elem> const& ptr) -> void;
0154    void cancel(std::shared_ptr<elem> const& ptr);
0155    auto reset() -> void;
0156 
0157    [[nodiscard]]
0158    auto const& get_parser() const noexcept
0159    {
0160       return parser_;
0161    }
0162 
0163    auto cancel_waiting() -> std::size_t;
0164 
0165    // To be called exactly once to clean up state after a connection becomes unhealthy.
0166    // Requests are canceled or returned to the waiting state to be re-sent to the server,
0167    // depending on their configuration. After this function is called, prepare_write,
0168    // commit_write and consume_next must not be called until a reset() happens.
0169    // Otherwise, race conditions like the following might happen
0170    // (see https://github.com/boostorg/redis/pull/309 and https://github.com/boostorg/redis/issues/181):
0171    //
0172    //   - This function runs and cancels a request, then consume_next runs. It tries to access
0173    //     a request and adapter that might have been destroyed.
0174    //   - This function runs and returns a request to waiting, then prepare_write runs.
0175    //     It incorrectly sets the request state to staged, causing de synchronization between requests and responses.
0176    void cancel_on_conn_lost();
0177 
0178    [[nodiscard]]
0179    auto get_write_buffer() const noexcept -> std::string_view
0180    {
0181       return std::string_view{write_buffer_}.substr(write_offset_);
0182    }
0183 
0184    [[nodiscard]]
0185    auto get_prepared_read_buffer() noexcept -> read_buffer::span_type;
0186 
0187    [[nodiscard]]
0188    auto prepare_read() noexcept -> system::error_code;
0189 
0190    void commit_read(std::size_t read_size);
0191 
0192    [[nodiscard]]
0193    auto get_read_buffer_size() const noexcept -> std::size_t;
0194 
0195    void set_receive_adapter(any_adapter adapter);
0196 
0197    [[nodiscard]]
0198    auto get_usage() const noexcept -> usage
0199    {
0200       return usage_;
0201    }
0202 
0203    void set_config(config const& cfg);
0204 
0205 private:
0206    void commit_usage(bool is_push, read_buffer::consume_result res);
0207 
0208    [[nodiscard]]
0209    auto is_next_push(std::string_view data) const noexcept -> bool;
0210 
0211    // Completes requests that don't expect a response
0212    void release_push_requests();
0213 
0214    [[nodiscard]]
0215    consume_result consume_impl(system::error_code& ec);
0216 
0217    read_buffer read_buffer_;
0218    std::string write_buffer_;
0219    std::size_t write_offset_{};  // how many bytes of the write buffer have been written?
0220    std::deque<std::shared_ptr<elem>> reqs_;
0221    resp3::parser parser_{};
0222    bool on_push_ = false;
0223    bool cancel_run_called_ = false;
0224    usage usage_;
0225    any_adapter receive_adapter_;
0226 };
0227 
0228 auto make_elem(request const& req, any_adapter adapter) -> std::shared_ptr<multiplexer::elem>;
0229 
0230 }  // namespace detail
0231 }  // namespace boost::redis
0232 
0233 #endif  // BOOST_REDIS_MULTIPLEXER_HPP