diff --git a/include/rest_rpc/md5.hpp b/include/rest_rpc/md5.hpp index fcdb1fc..e5812da 100644 --- a/include/rest_rpc/md5.hpp +++ b/include/rest_rpc/md5.hpp @@ -271,37 +271,23 @@ struct MD5CE { } ////////////////////////////////////////////////////////////////////////////// // HELPER FUNCTIONS - static constexpr uint32_t StringLength(const char *string) { - const char *end = string; - while (*end != 0) - ++end; - // Double check that the precision losing conversion is safe. - // DCHECK(end >= string); - // DCHECK(static_cast(static_cast(end - - // string)) == - // (end - string)); - return static_cast(end - string); - } static constexpr uint32_t SwapEndian(uint32_t a) { return ((a & 0xff) << 24) | (((a >> 8) & 0xff) << 16) | (((a >> 16) & 0xff) << 8) | ((a >> 24) & 0xff); } ////////////////////////////////////////////////////////////////////////////// // WRAPPER FUNCTIONS - static constexpr uint64_t Hash64(const char *data, uint32_t n) { - IntermediateData intermediate = ProcessMessage(data, n); - return (static_cast(SwapEndian(intermediate.a)) << 32) | - static_cast(SwapEndian(intermediate.b)); - } + // static constexpr uint64_t Hash64(const char *data, uint32_t n) { + // IntermediateData intermediate = ProcessMessage(data, n); + // return (static_cast(SwapEndian(intermediate.a)) << 32) | + // static_cast(SwapEndian(intermediate.b)); + // } static constexpr uint32_t Hash32(const char *data, uint32_t n) { IntermediateData intermediate = ProcessMessage(data, n); return SwapEndian(intermediate.a); } }; // https://chromium.googlesource.com/chromium/src/base/+/refs/heads/main/hash/md5_constexpr_internal.h -constexpr uint32_t MD5Hash32(const char *string) { - return MD5CE::Hash32(string, MD5CE::StringLength(string)); -} constexpr uint32_t MD5Hash32(const char *string, uint32_t length) { return MD5CE::Hash32(string, length); } diff --git a/include/rest_rpc/rpc_client.hpp b/include/rest_rpc/rpc_client.hpp index 5167177..15b8373 100644 --- a/include/rest_rpc/rpc_client.hpp +++ b/include/rest_rpc/rpc_client.hpp @@ -20,9 +20,7 @@ template struct call_result { R value; }; -template <> struct call_result { - rpc_errc ec; -}; +template <> struct call_result { rpc_errc ec; }; class rpc_client { public: @@ -57,11 +55,6 @@ public: } auto it = endpoints.begin(); - if (it == endpoints.end()) { - REST_LOG_ERROR << "resolve failed"; - co_return std::make_error_code(std::errc::bad_address); - } - auto endpoint = it->endpoint(); auto conn_r = co_await (watchdog(duration) || socket_->impl_.async_connect( @@ -137,10 +130,6 @@ public: header.function_id = topic_id; auto [it, r] = socket_->sub_ops_.emplace(topic_id, sub_operation{}); - if (!r) { - REST_LOG_ERROR << "subscribe duplicate topic"; - co_return call_result{rpc_errc::duplicate_topic}; - } std::tie(b, ret) = co_await ( asio::async_compose( diff --git a/tests/test_rest_rpc.cpp b/tests/test_rest_rpc.cpp index 81d77b3..e2c85b3 100644 --- a/tests/test_rest_rpc.cpp +++ b/tests/test_rest_rpc.cpp @@ -210,6 +210,11 @@ asio::awaitable test_router() { TEST_CASE("test router") { sync_wait(get_global_executor(), test_router()); } +TEST_CASE("test context pool") { + io_context_pool pool(0); + CHECK(pool.size() == 1); +} + TEST_CASE("test server start") { rpc_server server("127.0.0.1:9005"); server.register_handler(); @@ -261,6 +266,7 @@ TEST_CASE("test server start") { auto result0 = sync_wait(cl.get_executor(), cl.call_for(std::chrono::minutes(2), "test")); + REST_LOG_INFO << make_error_code(result0.ec).message(); CHECK(result0.ec == rpc_errc::ok); result0 = @@ -275,10 +281,18 @@ TEST_CASE("test server start") { auto result1 = sync_wait( cl.get_executor(), cl.call_for(std::chrono::minutes(2), "test")); CHECK(result1.ec == rpc_errc::ok); + cl.enable_tcp_no_delay(false); auto result2 = sync_wait(cl.get_executor(), cl.call(1)); CHECK(result2.ec == rpc_errc::ok); } + { + auto result = + sync_wait(cl.get_executor(), + cl.call_for(std::chrono::milliseconds(0), 1, 2)); + CHECK(result.ec == rpc_errc::request_timeout); + } + ec = server.async_start(); CHECK(!ec); ec = server.async_start(); @@ -286,6 +300,11 @@ TEST_CASE("test server start") { server.stop(); server.stop(); + { + auto result = sync_wait(cl.get_executor(), cl.call(1, 2)); + CHECK(result.ec != rpc_errc::ok); + } + ec = server.async_start(); CHECK(ec); std::cout << ec.message() << "\n"; @@ -312,7 +331,16 @@ TEST_CASE("test server start") { } TEST_CASE("test cross ending") { + rpc_server server("127.0.0.1:9004"); + server.register_handler(); + server.enable_cross_ending(true); + server.async_start(); + rpc_client client{}; + client.enable_cross_ending(true); + sync_wait(get_global_executor(), client.connect("127.0.0.1:9004")); + auto ret = sync_wait(client.get_executor(), client.call("test")); + CHECK(ret.value == "test"); } TEST_CASE("test pub sub") { @@ -330,6 +358,7 @@ TEST_CASE("test pub sub") { CHECK(result.ec == rpc_errc::ok); CHECK(result.value == "publish message"); } + promise.set_value(); }; @@ -358,6 +387,25 @@ TEST_CASE("test reconnect") { CHECK(!ec); ec = sync_wait(get_global_executor(), client.connect("127.0.0.1:9005")); CHECK(!ec); + + ec = sync_wait(get_global_executor(), client.connect("127.0.0.1:9006")); + CHECK(ec); + REST_LOG_INFO << ec.message(); + ec = sync_wait( + get_global_executor(), + client.connect("127.0.0.0:9006", std::chrono::milliseconds(200))); + CHECK(ec); + REST_LOG_INFO << ec.message(); + ec = + sync_wait(get_global_executor(), + client.connect("127.0.0.x:9006", std::chrono::milliseconds(0))); + CHECK(ec); + REST_LOG_INFO << ec.message(); + ec = sync_wait( + get_global_executor(), + client.connect("127.0.0.x:9006", std::chrono::milliseconds(200))); + CHECK(ec); + REST_LOG_INFO << ec.message(); } // doctest comments