mirror of
https://github.com/qicosmos/rest_rpc.git
synced 2026-08-29 08:34:47 +08:00
simplify and rename codec
This commit is contained in:
+19
-49
@@ -6,64 +6,21 @@
|
||||
#include <msgpack.hpp>
|
||||
|
||||
namespace rest_rpc {
|
||||
namespace rpc_service {
|
||||
|
||||
template <typename Arg> auto pack_one(Arg &&arg) {
|
||||
if constexpr (util::CharArrayRef<Arg> || util::CharArray<Arg> ||
|
||||
util::string<Arg>) {
|
||||
return std::string_view(std::forward<Arg>(arg));
|
||||
} else {
|
||||
return std::to_string(arg);
|
||||
}
|
||||
}
|
||||
|
||||
using buffer_type = msgpack::sbuffer;
|
||||
struct msgpack_codec {
|
||||
const static size_t init_size = 2 * 1024;
|
||||
|
||||
template <typename... Args> static auto pack_args(Args &&...args) {
|
||||
template <typename... Args> inline static auto pack_args(Args &&...args) {
|
||||
if constexpr (sizeof...(Args) == 0) {
|
||||
return std::string_view{};
|
||||
} else if constexpr (sizeof...(Args) == 1 && util::is_basic_v<Args...>) {
|
||||
return pack_one(std::forward<Args>(args)...);
|
||||
} else {
|
||||
buffer_type buffer(init_size);
|
||||
msgpack::sbuffer buffer(2 * 1024);
|
||||
msgpack::pack(buffer, std::forward_as_tuple(std::forward<Args>(args)...));
|
||||
return std::string(buffer.data(), buffer.size());
|
||||
}
|
||||
}
|
||||
|
||||
template <typename Arg> static std::string pack_to_string(Arg &&arg) {
|
||||
buffer_type buffer(init_size);
|
||||
msgpack::pack(buffer, arg);
|
||||
return std::string(buffer.data(), buffer.size());
|
||||
}
|
||||
|
||||
template <typename Arg, typename... Args,
|
||||
typename = typename std::enable_if<std::is_enum<Arg>::value>::type>
|
||||
static std::string pack_args_str(Arg arg, Args &&...args) {
|
||||
buffer_type buffer(init_size);
|
||||
msgpack::pack(buffer,
|
||||
std::forward_as_tuple((int)arg, std::forward<Args>(args)...));
|
||||
return std::string(buffer.data(), buffer.size());
|
||||
}
|
||||
|
||||
template <typename T> buffer_type pack(T &&t) const {
|
||||
buffer_type buffer;
|
||||
msgpack::pack(buffer, std::forward<T>(t));
|
||||
return buffer;
|
||||
}
|
||||
|
||||
template <typename T> T unpack(char const *data, size_t length) {
|
||||
try {
|
||||
msgpack::unpack(msg_, data, length);
|
||||
return msg_.get().as<T>();
|
||||
} catch (...) {
|
||||
throw std::invalid_argument("unpack failed: Args not match!");
|
||||
}
|
||||
}
|
||||
|
||||
template <typename T> T unpack(std::string_view data) {
|
||||
template <typename T> inline static T unpack(std::string_view data) {
|
||||
if constexpr (std::is_fundamental_v<T>) {
|
||||
T t;
|
||||
auto r = std::from_chars(data.data(), data.data() + data.size(), t);
|
||||
@@ -76,14 +33,27 @@ struct msgpack_codec {
|
||||
} else if constexpr (std::is_same_v<std::string_view, T>) {
|
||||
return data;
|
||||
} else {
|
||||
return unpack<T>(data.data(), data.size());
|
||||
try {
|
||||
static msgpack::unpacked msg;
|
||||
msgpack::unpack(msg, data.data(), data.size());
|
||||
return msg.get().as<T>();
|
||||
} catch (...) {
|
||||
throw std::invalid_argument("unpack failed: Args not match!");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private:
|
||||
msgpack::unpacked msg_;
|
||||
template <typename Arg> inline static auto pack_one(Arg &&arg) {
|
||||
if constexpr (util::CharArrayRef<Arg> || util::CharArray<Arg> ||
|
||||
util::string<Arg>) {
|
||||
return std::string_view(std::forward<Arg>(arg));
|
||||
} else {
|
||||
return std::to_string(arg);
|
||||
}
|
||||
}
|
||||
};
|
||||
} // namespace rpc_service
|
||||
|
||||
} // namespace rest_rpc
|
||||
|
||||
#endif // REST_RPC_CODEC_H_
|
||||
@@ -164,8 +164,7 @@ private:
|
||||
template <typename R, typename... Args>
|
||||
asio::awaitable<call_result<R>> call_impl(rest_rpc_header &header,
|
||||
Args &&...args) {
|
||||
rpc_service::msgpack_codec codec;
|
||||
auto buf = codec.pack_args(std::forward<Args>(args)...);
|
||||
auto buf = msgpack_codec::pack_args(std::forward<Args>(args)...);
|
||||
header.body_len = buf.size();
|
||||
if (cross_ending_) {
|
||||
prepare_for_send(header);
|
||||
@@ -227,12 +226,11 @@ private:
|
||||
}
|
||||
result.ec = (rpc_errc)socket_->body_[0];
|
||||
if constexpr (!std::is_void_v<R>) {
|
||||
rpc_service::msgpack_codec codec;
|
||||
if constexpr (util::is_basic_v<R>) {
|
||||
result.value = codec.unpack<R>(std::string_view(
|
||||
result.value = msgpack_codec::unpack<R>(std::string_view(
|
||||
socket_->body_.data() + 1, resp_header.body_len - 1));
|
||||
} else {
|
||||
auto tp = codec.unpack<std::tuple<R>>(std::string_view(
|
||||
auto tp = msgpack_codec::unpack<std::tuple<R>>(std::string_view(
|
||||
socket_->body_.data() + 1, resp_header.body_len - 1));
|
||||
result.value = std::move(std::get<0>(tp));
|
||||
}
|
||||
|
||||
@@ -237,8 +237,8 @@ asio::awaitable<std::error_code> rpc_context::response(Args &&...args) {
|
||||
REST_LOG_ERROR << "rpc context init failed";
|
||||
co_return make_error_code(rpc_errc::rpc_context_init_failed);
|
||||
}
|
||||
rpc_service::msgpack_codec codec;
|
||||
rpc_result result(codec.pack_args(std::forward<Args>(args)...));
|
||||
|
||||
rpc_result result(msgpack_codec::pack_args(std::forward<Args>(args)...));
|
||||
has_response_ = true;
|
||||
co_return co_await conn_->response(result);
|
||||
}
|
||||
|
||||
@@ -72,7 +72,6 @@ public:
|
||||
asio::awaitable<rpc_result> route(uint32_t key, std::string_view data) {
|
||||
rpc_result route_result{};
|
||||
try {
|
||||
rpc_service::msgpack_codec codec;
|
||||
auto it = map_invokers_.find(key);
|
||||
if (it == map_invokers_.end()) {
|
||||
route_result.result = "unknown function: " + get_name_by_key(key);
|
||||
@@ -149,15 +148,15 @@ private:
|
||||
} else {
|
||||
if constexpr (std::is_void_v<Self>) {
|
||||
if constexpr (is_awaitable_v<R>) {
|
||||
ret = rpc_service::msgpack_codec::pack_args(co_await f());
|
||||
ret = msgpack_codec::pack_args(co_await f());
|
||||
} else {
|
||||
ret = rpc_service::msgpack_codec::pack_args(f());
|
||||
ret = msgpack_codec::pack_args(f());
|
||||
}
|
||||
} else {
|
||||
if constexpr (is_awaitable_v<R>) {
|
||||
ret = rpc_service::msgpack_codec::pack_args(co_await (*self.*f)());
|
||||
ret = msgpack_codec::pack_args(co_await (*self.*f)());
|
||||
} else {
|
||||
ret = rpc_service::msgpack_codec::pack_args((*self.*f)());
|
||||
ret = msgpack_codec::pack_args((*self.*f)());
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -167,37 +166,35 @@ private:
|
||||
template <typename R, typename Arg, typename F, typename Self>
|
||||
asio::awaitable<void> handle_one_arg(std::string_view str, const F &f,
|
||||
rpc_result &ret, Self *self) {
|
||||
rpc_service::msgpack_codec codec;
|
||||
if constexpr (is_void_v<R>) {
|
||||
if constexpr (std::is_void_v<Self>) {
|
||||
if constexpr (is_awaitable_v<R>) {
|
||||
co_await f(codec.unpack<Arg>(str));
|
||||
co_await f(msgpack_codec::unpack<Arg>(str));
|
||||
} else {
|
||||
f(codec.unpack<Arg>(str));
|
||||
f(msgpack_codec::unpack<Arg>(str));
|
||||
}
|
||||
} else {
|
||||
if constexpr (is_awaitable_v<R>) {
|
||||
co_await (*self.*f)(codec.unpack<Arg>(str));
|
||||
co_await (*self.*f)(msgpack_codec::unpack<Arg>(str));
|
||||
} else {
|
||||
(*self.*f)(codec.unpack<Arg>(str));
|
||||
(*self.*f)(msgpack_codec::unpack<Arg>(str));
|
||||
}
|
||||
}
|
||||
} else {
|
||||
if constexpr (std::is_void_v<Self>) {
|
||||
if constexpr (is_awaitable_v<R>) {
|
||||
ret = rpc_service::msgpack_codec::pack_args(
|
||||
co_await f(codec.unpack<Arg>(str)));
|
||||
ret = msgpack_codec::pack_args(
|
||||
co_await f(msgpack_codec::unpack<Arg>(str)));
|
||||
} else {
|
||||
ret =
|
||||
rpc_service::msgpack_codec::pack_args(f(codec.unpack<Arg>(str)));
|
||||
ret = msgpack_codec::pack_args(f(msgpack_codec::unpack<Arg>(str)));
|
||||
}
|
||||
} else {
|
||||
if constexpr (is_awaitable_v<R>) {
|
||||
ret = rpc_service::msgpack_codec::pack_args(
|
||||
co_await (*self.*f)(codec.unpack<Arg>(str)));
|
||||
ret = msgpack_codec::pack_args(
|
||||
co_await (*self.*f)(msgpack_codec::unpack<Arg>(str)));
|
||||
} else {
|
||||
ret = rpc_service::msgpack_codec::pack_args(
|
||||
(*self.*f)(codec.unpack<Arg>(str)));
|
||||
ret = msgpack_codec::pack_args(
|
||||
(*self.*f)(msgpack_codec::unpack<Arg>(str)));
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -207,8 +204,7 @@ private:
|
||||
template <typename R, typename Args, typename F, typename Self>
|
||||
asio::awaitable<void> handle_more_args(std::string_view str, const F &f,
|
||||
rpc_result &ret, Self *self) {
|
||||
rpc_service::msgpack_codec codec;
|
||||
auto tp = codec.unpack<Args>(str);
|
||||
auto tp = msgpack_codec::unpack<Args>(str);
|
||||
if constexpr (std::is_void_v<R>) {
|
||||
if constexpr (std::is_void_v<Self>) {
|
||||
if constexpr (is_awaitable_v<R>) {
|
||||
@@ -234,20 +230,19 @@ private:
|
||||
} else {
|
||||
if constexpr (std::is_void_v<Self>) {
|
||||
if constexpr (is_awaitable_v<R>) {
|
||||
ret =
|
||||
rpc_service::msgpack_codec::pack_args(co_await std::apply(f, tp));
|
||||
ret = msgpack_codec::pack_args(co_await std::apply(f, tp));
|
||||
} else {
|
||||
ret = rpc_service::msgpack_codec::pack_args(std::apply(f, tp));
|
||||
ret = msgpack_codec::pack_args(std::apply(f, tp));
|
||||
}
|
||||
} else {
|
||||
if constexpr (is_awaitable_v<R>) {
|
||||
ret = rpc_service::msgpack_codec::pack_args(co_await std::apply(
|
||||
ret = msgpack_codec::pack_args(co_await std::apply(
|
||||
[self, &f](auto &&...args) {
|
||||
return (*self.*f)(std::forward<decltype(args)>(args)...);
|
||||
},
|
||||
tp));
|
||||
} else {
|
||||
ret = rpc_service::msgpack_codec::pack_args(std::apply(
|
||||
ret = msgpack_codec::pack_args(std::apply(
|
||||
[self, &f](auto &&...args) {
|
||||
return (*self.*f)(std::forward<decltype(args)>(args)...);
|
||||
},
|
||||
|
||||
@@ -113,8 +113,8 @@ public:
|
||||
auto conns = get_connections();
|
||||
for (auto &[_, conn] : conns) {
|
||||
if (conn->topic_id() == id) {
|
||||
co_await conn->response(
|
||||
rpc_service::msgpack_codec::pack_args(std::forward<T>(t)), id);
|
||||
co_await conn->response(msgpack_codec::pack_args(std::forward<T>(t)),
|
||||
id);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+14
-16
@@ -241,11 +241,10 @@ asio::awaitable<void> test_router() {
|
||||
auto ret2 = co_await router.route(get_key<&dummy::echo_coro>(), "test");
|
||||
CHECK(ret2.ec == rpc_errc::ok);
|
||||
}
|
||||
rpc_service::msgpack_codec codec;
|
||||
|
||||
{
|
||||
auto s = codec.pack_args(1);
|
||||
auto s1 = codec.pack_args("test");
|
||||
auto s = msgpack_codec::pack_args(1);
|
||||
auto s1 = msgpack_codec::pack_args("test");
|
||||
auto r = co_await router.route(get_key<round1>(), s);
|
||||
auto r1 = co_await router.route(get_key<echo>(), s1);
|
||||
|
||||
@@ -254,7 +253,7 @@ asio::awaitable<void> test_router() {
|
||||
std::cout << "\n";
|
||||
}
|
||||
|
||||
auto args = codec.pack_args(1, 2);
|
||||
auto args = msgpack_codec::pack_args(1, 2);
|
||||
std::string_view str(args.data(), args.size());
|
||||
|
||||
{
|
||||
@@ -263,12 +262,12 @@ asio::awaitable<void> test_router() {
|
||||
std::cout << "\n";
|
||||
}
|
||||
|
||||
auto args1 = codec.pack_args("it is a test");
|
||||
auto args1 = msgpack_codec::pack_args("it is a test");
|
||||
std::string_view str1(args1.data(), args1.size());
|
||||
|
||||
{
|
||||
auto result = co_await router.route(get_key<&dummy::add>(), str);
|
||||
auto r = codec.unpack<int>(result.result);
|
||||
auto r = msgpack_codec::unpack<int>(result.result);
|
||||
auto result1 = co_await router.route(get_key<&dummy::foo>(), str1);
|
||||
CHECK(r == 3);
|
||||
CHECK(result1.ec == rpc_errc::ok);
|
||||
@@ -291,14 +290,13 @@ TEST_CASE("test rpc_connection") {
|
||||
bool cross_ending_ = false;
|
||||
rpc_router router;
|
||||
router.register_handler<get_person>();
|
||||
rpc_service::msgpack_codec codec;
|
||||
|
||||
person p{1, "tom", 20};
|
||||
|
||||
auto buf = codec.pack_to_string(std::tuple(p));
|
||||
auto buf = msgpack_codec::pack_args(p);
|
||||
auto ret = sync_wait(get_global_executor(),
|
||||
router.route(get_key<get_person>(), buf));
|
||||
auto tp =
|
||||
codec.unpack<std::tuple<person>>(ret.data().data(), ret.data().size());
|
||||
auto tp = msgpack_codec::unpack<std::tuple<person>>(ret.data());
|
||||
dummy d{};
|
||||
router.register_handler<&dummy::add>(&d);
|
||||
auto conn = std::make_shared<rpc_connection>(std::move(socket), conn_id,
|
||||
@@ -336,12 +334,12 @@ TEST_CASE("test server start") {
|
||||
|
||||
static_assert(util::CharArrayRef<char const(&)[5]>);
|
||||
static_assert(util::CharArray<const char[5]>);
|
||||
rpc_service::msgpack_codec::pack_args();
|
||||
rpc_service::msgpack_codec::pack_args(1, 2);
|
||||
auto s1 = rpc_service::msgpack_codec::pack_args("test");
|
||||
auto s2 = rpc_service::msgpack_codec::pack_args(std::string_view("test2"));
|
||||
auto s3 = rpc_service::msgpack_codec::pack_args(std::string("test2"));
|
||||
auto s5 = rpc_service::msgpack_codec::pack_args(123);
|
||||
msgpack_codec::pack_args();
|
||||
msgpack_codec::pack_args(1, 2);
|
||||
auto s1 = msgpack_codec::pack_args("test");
|
||||
auto s2 = msgpack_codec::pack_args(std::string_view("test2"));
|
||||
auto s3 = msgpack_codec::pack_args(std::string("test2"));
|
||||
auto s5 = msgpack_codec::pack_args(123);
|
||||
|
||||
// auto future = asio::co_spawn(cl.get_executor(),
|
||||
// cl.connect("127.0.0.1:9005"), asio::use_future); auto conn_ec =
|
||||
|
||||
Reference in New Issue
Block a user