|
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