add unsafe response

This commit is contained in:
qicosmos
2025-10-04 16:54:07 +08:00
parent 500303ad66
commit d5b6edc907
3 changed files with 49 additions and 20 deletions
+1 -1
View File
@@ -6,7 +6,7 @@
#include <arpa/inet.h> #include <arpa/inet.h>
#endif #endif
#ifdef __APPLE__ #if defined(__APPLE__) || defined(_WIN32)
#else #else
inline uint64_t htonll(uint64_t value) { inline uint64_t htonll(uint64_t value) {
return ((uint64_t)htonl(value & 0xFFFFFFFF) << 32) | htonl(value >> 32); return ((uint64_t)htonl(value & 0xFFFFFFFF) << 32) | htonl(value >> 32);
+22 -5
View File
@@ -27,10 +27,15 @@ public:
auto get_executor(); auto get_executor();
template <auto func, typename... Args> template <auto func, typename... Args>
asio::awaitable<std::error_code> response(Args &&...args); asio::awaitable<std::error_code> response_s(Args &&...args);
template <auto func, typename... Args> template <auto func, typename... Args>
std::error_code sync_response(Args &&...args); std::error_code sync_response_s(Args &&...args);
template <typename... Args>
asio::awaitable<std::error_code> response(Args &&...args);
template <typename... Args> std::error_code sync_response(Args &&...args);
private: private:
rpc_context() = default; rpc_context() = default;
@@ -177,12 +182,24 @@ private:
// zero or one arguments // zero or one arguments
template <auto func, typename... Args> template <auto func, typename... Args>
asio::awaitable<std::error_code> rpc_context::response(Args &&...args) { asio::awaitable<std::error_code> rpc_context::response_s(Args &&...args) {
using args_tuple = using args_tuple =
typename util::function_traits<decltype(func)>::return_type; typename util::function_traits<decltype(func)>::return_type;
static_assert( static_assert(
std::is_constructible_v<args_tuple, Args...>, std::is_constructible_v<args_tuple, Args...>,
"rpc function return type and response arguments are not match"); "rpc function return type and response arguments are not match");
return response(std::forward<Args>(args)...);
}
template <auto func, typename... Args>
std::error_code rpc_context::sync_response_s(Args &&...args) {
return sync_wait(conn_->get_executor(),
response_s<func>(std::forward<Args>(args)...));
}
template <typename... Args>
asio::awaitable<std::error_code> rpc_context::response(Args &&...args) {
if (has_response_) { if (has_response_) {
co_return make_error_code(rpc_errc::has_response); co_return make_error_code(rpc_errc::has_response);
} }
@@ -193,10 +210,10 @@ asio::awaitable<std::error_code> rpc_context::response(Args &&...args) {
co_return co_await conn_->response(result); co_return co_await conn_->response(result);
} }
template <auto func, typename... Args> template <typename... Args>
std::error_code rpc_context::sync_response(Args &&...args) { std::error_code rpc_context::sync_response(Args &&...args) {
return sync_wait(conn_->get_executor(), return sync_wait(conn_->get_executor(),
response<func>(std::forward<Args>(args)...)); response(std::forward<Args>(args)...));
} }
auto rpc_context::get_executor() { return conn_->get_executor(); } auto rpc_context::get_executor() { return conn_->get_executor(); }
+26 -14
View File
@@ -42,11 +42,15 @@ std::string_view echo_sv(std::string_view str) { return str; }
template<auto func> template<auto func>
asio::awaitable<void> response(auto ctx) { asio::awaitable<void> response(auto ctx) {
auto ec = co_await ctx.template response<func>("test"); auto ec = co_await ctx.template response_s<func>("test");
if (ec) { if (ec) {
REST_LOG_ERROR << "response error: " << ec.message(); REST_LOG_ERROR << "response error: " << ec.message();
} }
ec = co_await ctx.template response<func>("test"); ec = co_await ctx.template response_s<func>("test");
REST_LOG_ERROR << ec.message();
CHECK(ec);
ec = co_await ctx.response("test");
REST_LOG_ERROR << ec.message(); REST_LOG_ERROR << ec.message();
CHECK(ec); CHECK(ec);
} }
@@ -56,17 +60,20 @@ std::string_view delay_response(std::string_view str) {
// set_delay before response in another thread // set_delay before response in another thread
ctx.set_delay(true); ctx.set_delay(true);
std::thread thd([ctx = std::move(ctx)]() mutable { // std::thread thd([ctx = std::move(ctx)]() mutable {
std::this_thread::sleep_for(std::chrono::seconds(2)); // std::this_thread::sleep_for(std::chrono::seconds(2));
// auto ec = ctx.sync_response("it is from a detached thread"); // // auto ec = ctx.sync_response("it is from a detached thread");
// if (ec) { // // if (ec) {
// REST_LOG_ERROR << "response error: " << ec.message(); // // REST_LOG_ERROR << "response error: " << ec.message();
// } // // }
auto executor = ctx.get_executor(); // auto executor = ctx.get_executor();
// async_start(executor, response<delay_response>(std::move(ctx))); // // async_start(executor, response<delay_response>(std::move(ctx)));
asio::co_spawn(executor, response<delay_response>(std::move(ctx)), asio::detached); // asio::co_spawn(executor, response<delay_response>(std::move(ctx)), asio::detached);
}); // });
thd.detach(); // thd.detach();
auto executor = ctx.get_executor();
asio::co_spawn(executor, response<delay_response>(std::move(ctx)), asio::detached);
// this return value is meaningless, because it will response later, the // this return value is meaningless, because it will response later, the
// return type is important for client, so just return an empty value here. // return type is important for client, so just return an empty value here.
@@ -89,7 +96,7 @@ asio::awaitable<std::string> no_arg_coro1() {
co_return "test"; co_return "test";
} }
// TODO: handle connection lifetime, client pool, pub/sub // TODO: client pool, pub/sub
asio::awaitable<void> test_router() { asio::awaitable<void> test_router() {
rpc_router router; rpc_router router;
router.register_handler<add>(); router.register_handler<add>();
@@ -232,6 +239,11 @@ TEST_CASE("test server start") {
cl.call_for<delay_response>(std::chrono::minutes(2), "test")); cl.call_for<delay_response>(std::chrono::minutes(2), "test"));
CHECK(result0.ec == rpc_errc::ok); CHECK(result0.ec == rpc_errc::ok);
result0 =
sync_wait(cl.get_executor(),
cl.call_for<delay_response>(std::chrono::minutes(2), "test"));
CHECK(result0.ec == rpc_errc::ok);
auto result = auto result =
sync_wait(cl.get_executor(), sync_wait(cl.get_executor(),
cl.call_for<echo_sv>(std::chrono::minutes(2), "test")); cl.call_for<echo_sv>(std::chrono::minutes(2), "test"));