This commit is contained in:
qicosmos
2025-10-08 20:07:43 +08:00
parent b194d923db
commit eb87b456d7
3 changed files with 54 additions and 31 deletions
+5 -19
View File
@@ -271,37 +271,23 @@ struct MD5CE {
} }
////////////////////////////////////////////////////////////////////////////// //////////////////////////////////////////////////////////////////////////////
// HELPER FUNCTIONS // 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<std::ptrdiff_t>(static_cast<uint32_t>(end -
// string)) ==
// (end - string));
return static_cast<uint32_t>(end - string);
}
static constexpr uint32_t SwapEndian(uint32_t a) { static constexpr uint32_t SwapEndian(uint32_t a) {
return ((a & 0xff) << 24) | (((a >> 8) & 0xff) << 16) | return ((a & 0xff) << 24) | (((a >> 8) & 0xff) << 16) |
(((a >> 16) & 0xff) << 8) | ((a >> 24) & 0xff); (((a >> 16) & 0xff) << 8) | ((a >> 24) & 0xff);
} }
////////////////////////////////////////////////////////////////////////////// //////////////////////////////////////////////////////////////////////////////
// WRAPPER FUNCTIONS // WRAPPER FUNCTIONS
static constexpr uint64_t Hash64(const char *data, uint32_t n) { // static constexpr uint64_t Hash64(const char *data, uint32_t n) {
IntermediateData intermediate = ProcessMessage(data, n); // IntermediateData intermediate = ProcessMessage(data, n);
return (static_cast<uint64_t>(SwapEndian(intermediate.a)) << 32) | // return (static_cast<uint64_t>(SwapEndian(intermediate.a)) << 32) |
static_cast<uint64_t>(SwapEndian(intermediate.b)); // static_cast<uint64_t>(SwapEndian(intermediate.b));
} // }
static constexpr uint32_t Hash32(const char *data, uint32_t n) { static constexpr uint32_t Hash32(const char *data, uint32_t n) {
IntermediateData intermediate = ProcessMessage(data, n); IntermediateData intermediate = ProcessMessage(data, n);
return SwapEndian(intermediate.a); return SwapEndian(intermediate.a);
} }
}; };
// https://chromium.googlesource.com/chromium/src/base/+/refs/heads/main/hash/md5_constexpr_internal.h // 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) { constexpr uint32_t MD5Hash32(const char *string, uint32_t length) {
return MD5CE::Hash32(string, length); return MD5CE::Hash32(string, length);
} }
+1 -12
View File
@@ -20,9 +20,7 @@ template <typename R> struct call_result {
R value; R value;
}; };
template <> struct call_result<void> { template <> struct call_result<void> { rpc_errc ec; };
rpc_errc ec;
};
class rpc_client { class rpc_client {
public: public:
@@ -57,11 +55,6 @@ public:
} }
auto it = endpoints.begin(); 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 endpoint = it->endpoint();
auto conn_r = co_await (watchdog(duration) || auto conn_r = co_await (watchdog(duration) ||
socket_->impl_.async_connect( socket_->impl_.async_connect(
@@ -137,10 +130,6 @@ public:
header.function_id = topic_id; header.function_id = topic_id;
auto [it, r] = socket_->sub_ops_.emplace(topic_id, sub_operation{}); auto [it, r] = socket_->sub_ops_.emplace(topic_id, sub_operation{});
if (!r) {
REST_LOG_ERROR << "subscribe duplicate topic";
co_return call_result<R>{rpc_errc::duplicate_topic};
}
std::tie(b, ret) = co_await ( std::tie(b, ret) = co_await (
asio::async_compose<decltype(asio::use_awaitable), void(bool)>( asio::async_compose<decltype(asio::use_awaitable), void(bool)>(
+48
View File
@@ -210,6 +210,11 @@ asio::awaitable<void> test_router() {
TEST_CASE("test router") { sync_wait(get_global_executor(), 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") { TEST_CASE("test server start") {
rpc_server server("127.0.0.1:9005"); rpc_server server("127.0.0.1:9005");
server.register_handler<add>(); server.register_handler<add>();
@@ -261,6 +266,7 @@ TEST_CASE("test server start") {
auto result0 = auto result0 =
sync_wait(cl.get_executor(), sync_wait(cl.get_executor(),
cl.call_for<delay_response>(std::chrono::minutes(2), "test")); cl.call_for<delay_response>(std::chrono::minutes(2), "test"));
REST_LOG_INFO << make_error_code(result0.ec).message();
CHECK(result0.ec == rpc_errc::ok); CHECK(result0.ec == rpc_errc::ok);
result0 = result0 =
@@ -275,10 +281,18 @@ TEST_CASE("test server start") {
auto result1 = sync_wait( auto result1 = sync_wait(
cl.get_executor(), cl.call_for<echo>(std::chrono::minutes(2), "test")); cl.get_executor(), cl.call_for<echo>(std::chrono::minutes(2), "test"));
CHECK(result1.ec == rpc_errc::ok); CHECK(result1.ec == rpc_errc::ok);
cl.enable_tcp_no_delay(false);
auto result2 = sync_wait(cl.get_executor(), cl.call<round1>(1)); auto result2 = sync_wait(cl.get_executor(), cl.call<round1>(1));
CHECK(result2.ec == rpc_errc::ok); CHECK(result2.ec == rpc_errc::ok);
} }
{
auto result =
sync_wait(cl.get_executor(),
cl.call_for<add>(std::chrono::milliseconds(0), 1, 2));
CHECK(result.ec == rpc_errc::request_timeout);
}
ec = server.async_start(); ec = server.async_start();
CHECK(!ec); CHECK(!ec);
ec = server.async_start(); ec = server.async_start();
@@ -286,6 +300,11 @@ TEST_CASE("test server start") {
server.stop(); server.stop();
server.stop(); server.stop();
{
auto result = sync_wait(cl.get_executor(), cl.call<add>(1, 2));
CHECK(result.ec != rpc_errc::ok);
}
ec = server.async_start(); ec = server.async_start();
CHECK(ec); CHECK(ec);
std::cout << ec.message() << "\n"; std::cout << ec.message() << "\n";
@@ -312,7 +331,16 @@ TEST_CASE("test server start") {
} }
TEST_CASE("test cross ending") { TEST_CASE("test cross ending") {
rpc_server server("127.0.0.1:9004");
server.register_handler<echo>();
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<echo>("test"));
CHECK(ret.value == "test");
} }
TEST_CASE("test pub sub") { TEST_CASE("test pub sub") {
@@ -330,6 +358,7 @@ TEST_CASE("test pub sub") {
CHECK(result.ec == rpc_errc::ok); CHECK(result.ec == rpc_errc::ok);
CHECK(result.value == "publish message"); CHECK(result.value == "publish message");
} }
promise.set_value(); promise.set_value();
}; };
@@ -358,6 +387,25 @@ TEST_CASE("test reconnect") {
CHECK(!ec); CHECK(!ec);
ec = sync_wait(get_global_executor(), client.connect("127.0.0.1:9005")); ec = sync_wait(get_global_executor(), client.connect("127.0.0.1:9005"));
CHECK(!ec); 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 // doctest comments