From e1662e3a9fc289e333424d17eab596abf4acbbd6 Mon Sep 17 00:00:00 2001 From: "chufeng.qy@alibaba-inc.com" Date: Fri, 23 Dec 2022 10:44:20 +0800 Subject: [PATCH] add magic --- include/rest_rpc/connection.h | 11 ++++++----- include/rest_rpc/const_vars.h | 8 ++++---- include/rest_rpc/rpc_client.hpp | 17 +++++++++-------- 3 files changed, 19 insertions(+), 17 deletions(-) diff --git a/include/rest_rpc/connection.h b/include/rest_rpc/connection.h index 25e5120..6442507 100644 --- a/include/rest_rpc/connection.h +++ b/include/rest_rpc/connection.h @@ -228,11 +228,10 @@ private: void write() { auto &msg = write_queue_.front(); write_size_ = (uint32_t)msg.content->size(); - std::array write_buffers; - write_buffers[0] = asio::buffer(&write_size_, sizeof(uint32_t)); - write_buffers[1] = asio::buffer(&msg.req_id, sizeof(uint64_t)); - write_buffers[2] = asio::buffer(&msg.req_type, sizeof(request_type)); - write_buffers[3] = asio::buffer(msg.content->data(), write_size_); + header_ = {MAGIC_NUM, msg.req_type, write_size_, msg.req_id}; + std::array write_buffers; + write_buffers[0] = asio::buffer(&header_, sizeof(rpc_header)); + write_buffers[1] = asio::buffer(msg.content->data(), write_size_); auto self = this->shared_from_this(); async_write(write_buffers, @@ -397,6 +396,8 @@ private: std::uint64_t req_id_; request_type req_type_; + rpc_header header_; + uint32_t write_size_ = 0; std::mutex write_mtx_; diff --git a/include/rest_rpc/const_vars.h b/include/rest_rpc/const_vars.h index 8b3ba7c..bc68522 100644 --- a/include/rest_rpc/const_vars.h +++ b/include/rest_rpc/const_vars.h @@ -25,15 +25,15 @@ struct message_type { std::shared_ptr content; }; -#pragma pack(1) +static const uint8_t MAGIC_NUM = 39; struct rpc_header { + uint8_t magic; + request_type req_type; uint32_t body_len; uint64_t req_id; - request_type req_type; }; -#pragma pack() static const size_t MAX_BUF_LEN = 1048576 * 10; -static const size_t HEAD_LEN = 13; +static const size_t HEAD_LEN = sizeof(rpc_header); static const size_t INIT_BUF_SIZE = 2 * 1024; } // namespace rest_rpc \ No newline at end of file diff --git a/include/rest_rpc/rpc_client.hpp b/include/rest_rpc/rpc_client.hpp index fc375d5..8756d9e 100644 --- a/include/rest_rpc/rpc_client.hpp +++ b/include/rest_rpc/rpc_client.hpp @@ -445,11 +445,10 @@ private: void write() { auto &msg = outbox_[0]; write_size_ = (uint32_t)msg.content.length(); - std::array write_buffers; - write_buffers[0] = asio::buffer(&write_size_, sizeof(int32_t)); - write_buffers[1] = asio::buffer(&msg.req_id, sizeof(uint64_t)); - write_buffers[2] = asio::buffer(&msg.req_type, sizeof(request_type)); - write_buffers[3] = asio::buffer((char *)msg.content.data(), write_size_); + std::array write_buffers; + header_ = {MAGIC_NUM, msg.req_type, write_size_, msg.req_id}; + write_buffers[0] = asio::buffer(&header_, sizeof(rpc_header)); + write_buffers[1] = asio::buffer((char *)msg.content.data(), write_size_); async_write(write_buffers, [this](const asio::error_code &ec, const size_t length) { @@ -508,14 +507,14 @@ private: } } else { std::cout << ec.message() << "\n"; - + { std::unique_lock lock(cb_mtx_); - for (auto& item : callback_map_) { + for (auto &item : callback_map_) { item.second->callback(ec, {}); } } - + close(false); error_callback(ec); } @@ -854,6 +853,8 @@ private: char head_[HEAD_LEN] = {}; std::vector body_; + rpc_header header_; + std::unordered_map> sub_map_; std::set> key_token_set_;