|
6 | 6 | #include <tuple> |
7 | 7 |
|
8 | 8 | #include <parser/base_resp_parser.h> |
9 | | -#include <network/tcp_socket.hpp> |
10 | | -#include <network/unix_socket.hpp> |
| 9 | +#include <network/async_socket.hpp> |
11 | 10 |
|
12 | 11 | namespace async_redis |
13 | 12 | { |
14 | 13 | class connection |
15 | 14 | { |
16 | 15 | using async_socket = network::async_socket; |
17 | | - using tcp_socket = network::tcp_socket; |
18 | | - using unix_socket = network::unix_socket; |
19 | 16 |
|
20 | 17 | public: |
21 | 18 | using parser_t = parser::base_resp_parser::parser; |
22 | 19 | using reply_cb_t = std::function<void (parser_t)>; |
23 | 20 |
|
24 | | - connection(event_loop::event_loop_ev& event_loop) |
25 | | - : event_loop_(event_loop) { |
26 | | - } |
| 21 | + connection(event_loop::event_loop_ev& event_loop); |
27 | 22 |
|
28 | | - void connect(async_socket::connect_handler_t handler, const std::string& ip, int port) |
29 | | - { |
30 | | - if (!socket_ || !socket_->is_valid()) |
31 | | - socket_ = std::make_unique<tcp_socket>(event_loop_); |
| 23 | + void connect(async_socket::connect_handler_t handler, const std::string& ip, int port); |
| 24 | + void connect(async_socket::connect_handler_t handler, const std::string& path); |
32 | 25 |
|
33 | | - static_cast<tcp_socket*>(socket_.get())->async_connect(ip, port, handler); |
34 | | - } |
35 | | - |
36 | | - void connect(async_socket::connect_handler_t handler, const std::string& path) |
37 | | - { |
38 | | - if (!socket_ || !socket_->is_valid()) |
39 | | - socket_ = std::make_unique<unix_socket>(event_loop_); |
40 | | - |
41 | | - static_cast<unix_socket*>(socket_.get())->async_connect(path, handler); |
42 | | - } |
43 | | - |
44 | | - bool is_connected() const |
45 | | - { return socket_ && socket_->is_connected(); } |
46 | | - |
47 | | - inline int pressure() const |
48 | | - { return req_queue_.size(); } |
49 | | - |
50 | | - void disconnect() { |
51 | | - socket_->close(); |
52 | | - //TODO: check the policy! Should we free queue or retry again? |
53 | | - decltype(req_queue_) free_me; |
54 | | - free_me.swap(req_queue_); |
55 | | - } |
56 | | - |
57 | | - bool pipelined_send(std::string&& pipelined_cmds, std::vector<reply_cb_t>&& callbacks) |
58 | | - { |
59 | | - if (!is_connected()) |
60 | | - return false; |
61 | | - |
62 | | - return |
63 | | - socket_->async_write(pipelined_cmds, [this, cbs = std::move(callbacks)](ssize_t sent_chunk_len) { |
64 | | - if (sent_chunk_len == 0) |
65 | | - return disconnect(); |
66 | | - |
67 | | - if (!req_queue_.size() && cbs.size()) |
68 | | - do_read(); |
69 | | - |
70 | | - for(auto &&cb : cbs) |
71 | | - req_queue_.emplace(std::move(cb), nullptr); |
72 | | - }); |
73 | | - } |
74 | | - |
75 | | - bool send(const std::string&& command, const reply_cb_t& reply_cb) |
76 | | - { |
77 | | - if (!is_connected()) |
78 | | - return false; |
79 | | - |
80 | | - bool read_it = !req_queue_.size(); |
81 | | - req_queue_.emplace(reply_cb, nullptr); |
82 | | - |
83 | | - return |
84 | | - socket_->async_write(std::move(command), [this, read_it](ssize_t sent_chunk_len) { |
85 | | - if (sent_chunk_len == 0) |
86 | | - return disconnect(); |
87 | | - |
88 | | - if (read_it) |
89 | | - do_read(); |
90 | | - }); |
91 | | - } |
| 26 | + bool is_connected() const; |
| 27 | + void disconnect(); |
| 28 | + bool pipelined_send(std::string&& pipelined_cmds, std::vector<reply_cb_t>&& callbacks); |
| 29 | + bool send(const std::string&& command, const reply_cb_t& reply_cb); |
92 | 30 |
|
93 | 31 | private: |
94 | | - inline |
95 | | - void do_read() { |
96 | | - socket_->async_read(data_, max_data_size, std::bind(&connection::reply_received, this, std::placeholders::_1)); |
97 | | - } |
98 | | - |
99 | | - void reply_received(ssize_t len) |
100 | | - { |
101 | | - if (len == 0) |
102 | | - return disconnect(); |
103 | | - |
104 | | - ssize_t acc = 0; |
105 | | - while (acc < len && req_queue_.size()) |
106 | | - { |
107 | | - auto& request = req_queue_.front(); |
108 | | - |
109 | | - auto &cb = std::get<0>(request); |
110 | | - auto &parser = std::get<1>(request); |
111 | | - |
112 | | - bool is_finished = false; |
113 | | - acc += parser::base_resp_parser::append_chunk(parser, data_ + acc, len - acc, is_finished); |
114 | | - |
115 | | - if (!is_finished) |
116 | | - break; |
117 | | - |
118 | | - cb(parser); |
119 | | - req_queue_.pop(); //free the resources |
120 | | - } |
121 | | - |
122 | | - if (req_queue_.size()) |
123 | | - do_read(); |
124 | | - } |
| 32 | + void do_read(); |
| 33 | + void reply_received(ssize_t len); |
125 | 34 |
|
126 | 35 | private: |
127 | 36 | std::unique_ptr<async_socket> socket_; |
|
0 commit comments