mirror of
https://github.com/qicosmos/rest_rpc.git
synced 2026-08-29 08:34:47 +08:00
add magic
This commit is contained in:
@@ -228,11 +228,10 @@ private:
|
||||
void write() {
|
||||
auto &msg = write_queue_.front();
|
||||
write_size_ = (uint32_t)msg.content->size();
|
||||
std::array<asio::const_buffer, 4> 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<asio::const_buffer, 2> 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_;
|
||||
|
||||
|
||||
@@ -25,15 +25,15 @@ struct message_type {
|
||||
std::shared_ptr<std::string> 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
|
||||
@@ -445,11 +445,10 @@ private:
|
||||
void write() {
|
||||
auto &msg = outbox_[0];
|
||||
write_size_ = (uint32_t)msg.content.length();
|
||||
std::array<asio::const_buffer, 4> 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<asio::const_buffer, 2> 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<std::mutex> 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<char> body_;
|
||||
|
||||
rpc_header header_;
|
||||
|
||||
std::unordered_map<std::string, std::function<void(string_view)>> sub_map_;
|
||||
std::set<std::pair<std::string, std::string>> key_token_set_;
|
||||
|
||||
|
||||
Reference in New Issue
Block a user