From 16bf90cc43d4131cb3760fbb01fa452ca58c311f Mon Sep 17 00:00:00 2001 From: Ian Torres Date: Sun, 4 Jan 2026 13:37:36 -0300 Subject: [PATCH 1/4] Server test stabilization started. --- src/objects/session.cpp | 6 +- tests/server_base.hpp | 94 +++++++--- tests/server_test.cc | 378 ++++++++++++++++++++-------------------- 3 files changed, 263 insertions(+), 215 deletions(-) diff --git a/src/objects/session.cpp b/src/objects/session.cpp index 9c5137e..65b7d50 100644 --- a/src/objects/session.cpp +++ b/src/objects/session.cpp @@ -128,9 +128,9 @@ namespace aewt { {"action", "register"}, { "params", { - {"registered", _config->registered_ }, - {"clients_port", _config->clients_port_}, - {"sessions_port", _config->sessions_port_}, + {"registered", _config->registered_.load(std::memory_order_acquire) }, + {"clients_port", _config->clients_port_.load(std::memory_order_acquire) }, + {"sessions_port", _config->sessions_port_.load(std::memory_order_acquire)}, } }, }; diff --git a/tests/server_base.hpp b/tests/server_base.hpp index 5d58d88..f409c8a 100644 --- a/tests/server_base.hpp +++ b/tests/server_base.hpp @@ -20,64 +20,108 @@ #include class server_test : public testing::Test { - std::unique_ptr _remote_thread; - std::unique_ptr _local_thread; + std::unique_ptr thread_a_; + std::unique_ptr thread_b_; + std::unique_ptr thread_c_; protected: - std::shared_ptr _remote_server = std::make_shared(); - std::shared_ptr _local_server = std::make_shared(); + std::shared_ptr server_a_ = std::make_shared(); + std::shared_ptr server_b_ = std::make_shared(); + std::shared_ptr server_c_ = std::make_shared(); void SetUp() override { - _remote_thread = std::make_unique([this]() { - const auto &_config = _remote_server->get_config(); + thread_a_ = std::make_unique([this]() { + const auto &_config = server_a_->get_config(); _config->sessions_port_.store(0, std::memory_order_release); _config->clients_port_.store(0, std::memory_order_release); _config->repl_enabled = false; _config->threads_ = 4; - _remote_server->start(); + + LOG_INFO("starting server A"); + server_a_->start(); + LOG_INFO("server A stopped"); + }); + + LOG_INFO("waiting for server A ready"); + std::this_thread::sleep_for(std::chrono::seconds(3)); + + thread_b_ = std::make_unique([this]() { + const auto &_config = server_b_->get_config(); + _config->sessions_port_.store(0, std::memory_order_release); + _config->clients_port_.store(0, std::memory_order_release); + _config->is_node_ = true; + _config->threads_ = 4; + _config->repl_enabled = false; + + while (server_a_->get_config()->sessions_port_.load(std::memory_order_acquire) == 0 || + server_a_->get_config()->clients_port_.load(std::memory_order_acquire) == 0) { + std::this_thread::sleep_for(std::chrono::seconds(1)); + LOG_INFO("waiting for server A ready ..."); + } + + _config->remote_clients_port_.store( + server_a_->get_config()->clients_port_.load(std::memory_order_acquire), std::memory_order_release); + _config->remote_sessions_port_.store( + server_a_->get_config()->sessions_port_.load(std::memory_order_acquire), + std::memory_order_release); + + LOG_INFO("starting server B"); + server_b_->start(); + LOG_INFO("server B stopped"); }); + LOG_INFO("waiting for server B ready"); std::this_thread::sleep_for(std::chrono::seconds(5)); - _local_thread = std::make_unique([this]() { - const auto &_config = _local_server->get_config(); + thread_c_ = std::make_unique([this]() { + const auto &_config = server_c_->get_config(); _config->sessions_port_.store(0, std::memory_order_release); _config->clients_port_.store(0, std::memory_order_release); _config->is_node_ = true; _config->threads_ = 4; _config->repl_enabled = false; - while (_remote_server->get_config()->sessions_port_.load(std::memory_order_acquire) == 0 || _remote_server-> - get_config()->clients_port_.load(std::memory_order_acquire) == 0) { + while (server_a_->get_config()->sessions_port_.load(std::memory_order_acquire) == 0 || + server_a_->get_config()->clients_port_.load(std::memory_order_acquire) == 0) { std::this_thread::sleep_for(std::chrono::seconds(1)); - LOG_INFO("Waiting for remote ready ..."); + LOG_INFO("waiting for Server A ready ..."); } _config->remote_clients_port_.store( - _remote_server->get_config()->clients_port_.load(std::memory_order_acquire), std::memory_order_release); + server_a_->get_config()->clients_port_.load(std::memory_order_acquire), std::memory_order_release); _config->remote_sessions_port_.store( - _remote_server->get_config()->sessions_port_.load(std::memory_order_acquire), + server_a_->get_config()->sessions_port_.load(std::memory_order_acquire), std::memory_order_release); - _local_server->start(); + LOG_INFO("starting server C"); + server_c_->start(); + LOG_INFO("server C stopped"); }); - while (_local_server->get_state()->get_sessions().size() == 0 || _remote_server->get_state()->get_sessions().size() == 0) { - LOG_INFO("Waiting 1 second for remote and local ready ..."); + std::this_thread::sleep_for(std::chrono::seconds(5)); + + while (server_a_->get_state()->get_sessions().size() != 2 || !server_b_->get_config()->registered_.load(std::memory_order_acquire) || !server_c_->get_config()->registered_.load(std::memory_order_acquire)) { + LOG_INFO("waiting 1 second for all servers ready ..."); std::this_thread::sleep_for(std::chrono::seconds(1)); } } void TearDown() override { - LOG_INFO("Local server stop ..."); - _local_server->stop(); - LOG_INFO("Remote server stop ..."); - _remote_server->stop(); - while (!_local_server->get_state()->get_ioc().stopped() || !_remote_server->get_state()->get_ioc().stopped()) { - LOG_INFO("Waiting for io stop ..."); + LOG_INFO("server A stop ..."); + server_a_->stop(); + + LOG_INFO("server B stop ..."); + server_b_->stop(); + + LOG_INFO("server C stop ..."); + server_c_->stop(); + + while (!server_a_->get_state()->get_ioc().stopped() || !server_b_->get_state()->get_ioc().stopped() || !server_c_->get_state()->get_ioc().stopped()) { + LOG_INFO("waiting for io stop ..."); std::this_thread::sleep_for(std::chrono::seconds(1)); } - _remote_server.reset(); - _local_server.reset(); + server_a_.reset(); + server_b_.reset(); + server_c_.reset(); } }; diff --git a/tests/server_test.cc b/tests/server_test.cc index b6d523d..bd7e489 100644 --- a/tests/server_test.cc +++ b/tests/server_test.cc @@ -21,192 +21,196 @@ #include #include -TEST_F(server_test, assert_local_server_is_registered) { - ASSERT_TRUE(_local_server->get_config()->registered_); +TEST_F(server_test, assert_nodes_are_registered) { + ASSERT_TRUE(server_b_->get_config()->registered_); + ASSERT_TRUE(server_c_->get_config()->registered_); + ASSERT_TRUE(server_a_->get_state()->get_sessions().size() == 2); + ASSERT_TRUE(server_b_->get_state()->get_sessions().size() == 2); + ASSERT_TRUE(server_c_->get_state()->get_sessions().size() == 2); } -TEST_F(server_test, assert_local_server_accept_clients) { - boost::asio::io_context _ioc; - boost::asio::ip::tcp::resolver _resolver{make_strand(_ioc)}; - boost::beast::websocket::stream _client{make_strand(_ioc)}; - - auto const _results = _resolver.resolve("127.0.0.1", std::to_string(_local_server->get_config()->clients_port_.load(std::memory_order_acquire))); - boost::asio::connect(_client.next_layer(), _results); - - const auto _host = fmt::format("127.0.0.1:{}", std::to_string(_local_server->get_config()->clients_port_.load(std::memory_order_acquire))); - _client.handshake(_host, "/"); - - boost::beast::flat_buffer _accepted_buffer; - _client.read(_accepted_buffer); - std::cout << boost::beast::make_printable(_accepted_buffer.data()) << std::endl; - - ASSERT_TRUE(_local_server->get_state()->get_clients().size() == 1); - - boost::system::error_code ec; - _client.close(boost::beast::websocket::close_code::normal, ec); - _client.next_layer().close(ec); - - // Wait for processing. - std::this_thread::sleep_for(std::chrono::seconds(3)); - - ASSERT_TRUE(_local_server->get_state()->get_clients().size() == 0); -} - -TEST_F(server_test, assert_local_server_can_handle_subscribe) { - boost::asio::io_context _ioc; - boost::asio::ip::tcp::resolver _resolver{make_strand(_ioc)}; - boost::beast::websocket::stream _client{make_strand(_ioc)}; - - auto const _results = _resolver.resolve("127.0.0.1", std::to_string(_local_server->get_config()->clients_port_.load(std::memory_order_acquire))); - boost::asio::connect(_client.next_layer(), _results); - - const auto _host = fmt::format("127.0.0.1:{}", std::to_string(_local_server->get_config()->clients_port_.load(std::memory_order_acquire))); - _client.handshake(_host, "/"); - - boost::beast::flat_buffer _accepted_buffer; - _client.read(_accepted_buffer); - std::cout << boost::beast::make_printable(_accepted_buffer.data()) << std::endl; - - ASSERT_TRUE(_local_server->get_state()->get_clients().size() == 1); - - _client.write(boost::asio::buffer(std::string(serialize(boost::json::object{ - {"transaction_id", to_string(boost::uuids::random_generator()())}, - {"action", "subscribe"}, - {"params", {{"channel", "welcome"}}}, - })))); - - boost::beast::flat_buffer _buffer; - - _client.read(_buffer); - - std::cout << boost::beast::make_printable(_buffer.data()) << std::endl; - - boost::system::error_code ec; - _client.close(boost::beast::websocket::close_code::normal, ec); - _client.next_layer().close(ec); -} - -TEST_F(server_test, assert_local_server_can_handle_publish) { - boost::asio::io_context _ioc; - boost::asio::ip::tcp::resolver _resolver{make_strand(_ioc)}; - boost::beast::websocket::stream _local_client{make_strand(_ioc)}; - boost::beast::websocket::stream _other_local_client{make_strand(_ioc)}; - boost::beast::websocket::stream _remote_client{make_strand(_ioc)}; - boost::beast::websocket::stream _other_remote_client{make_strand(_ioc)}; - - { - auto const _results = _resolver.resolve("127.0.0.1", std::to_string(_local_server->get_config()->clients_port_.load(std::memory_order_acquire))); - boost::asio::connect(_local_client.next_layer(), _results); - boost::asio::connect(_other_local_client.next_layer(), _results); - } - - { - auto const _results = _resolver.resolve("127.0.0.1", std::to_string(_remote_server->get_config()->clients_port_.load(std::memory_order_acquire))); - boost::asio::connect(_remote_client.next_layer(), _results); - boost::asio::connect(_other_remote_client.next_layer(), _results); - } - - const auto _host = fmt::format("127.0.0.1:{}", std::to_string(_local_server->get_config()->clients_port_.load(std::memory_order_acquire))); - _local_client.handshake(_host, "/"); - _other_local_client.handshake(_host, "/"); - - _remote_client.handshake(_host, "/"); - _other_remote_client.handshake(_host, "/"); - - { - boost::beast::flat_buffer _buffer; - _local_client.read(_buffer); - std::cout << "Local client should receive accepted ..." << std::endl; - std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; - } - - { - boost::beast::flat_buffer _buffer; - _other_local_client.read(_buffer); - std::cout << "Other local client should receive accepted ..." << std::endl; - std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; - } - - { - boost::beast::flat_buffer _buffer; - _remote_client.read(_buffer); - std::cout << "Remote client should receive accepted ..." << std::endl; - std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; - } - - { - boost::beast::flat_buffer _buffer; - _other_remote_client.read(_buffer); - std::cout << "Other remote client should receive accepted ..." << std::endl; - std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; - _buffer.clear(); - } - - _local_client.write(boost::asio::buffer(std::string(serialize(boost::json::object{ - {"transaction_id", to_string(boost::uuids::random_generator()())}, - {"action", "subscribe"}, - {"params", {{"channel", "welcome"}}}, - })))); - - { - boost::beast::flat_buffer _buffer; - _local_client.read(_buffer); - std::cout << "Local client should receive subscribe ack ..." << std::endl; - std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; - _buffer.clear(); - } - - _remote_client.write(boost::asio::buffer(std::string(serialize(boost::json::object{ - {"transaction_id", to_string(boost::uuids::random_generator()())}, - {"action", "subscribe"}, - {"params", {{"channel", "welcome"}}}, - })))); - - { - boost::beast::flat_buffer _buffer; - _remote_client.read(_buffer); - std::cout << "Remote client should receive subscribe ack ..." << std::endl; - std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; - _buffer.clear(); - } - - std::this_thread::sleep_for(std::chrono::seconds(1)); - - _other_local_client.write(boost::asio::buffer(std::string(serialize(boost::json::object{ - {"transaction_id", to_string(boost::uuids::random_generator()())}, - {"action", "publish"}, - {"params", {{"channel", "welcome"}, {"payload", {{"message", "EHLO"}}}}}, - })))); - - { - boost::beast::flat_buffer _buffer; - _other_local_client.read(_buffer); - std::cout << "Other local client should receive publish ack ..." << std::endl; - std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; - _buffer.clear(); - } - - std::this_thread::sleep_for(std::chrono::seconds(1)); - - { - boost::beast::flat_buffer _buffer; - _local_client.read(_buffer); - std::cout << "Local client should receive publish message ..." << std::endl; - std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; - _buffer.clear(); - } - - - { - boost::beast::flat_buffer _buffer; - _remote_client.read(_buffer); - std::cout << "Remote client should receive publish message ..." << std::endl; - std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; - } - - boost::system::error_code ec; - _local_client.close(boost::beast::websocket::close_code::normal, ec); - _other_local_client.close(boost::beast::websocket::close_code::normal, ec); - _remote_client.close(boost::beast::websocket::close_code::normal, ec); - _other_remote_client.close(boost::beast::websocket::close_code::normal, ec); -} +// TEST_F(server_test, assert_local_server_accept_clients) { +// boost::asio::io_context _ioc; +// boost::asio::ip::tcp::resolver _resolver{make_strand(_ioc)}; +// boost::beast::websocket::stream _client{make_strand(_ioc)}; +// +// auto const _results = _resolver.resolve("127.0.0.1", std::to_string(_local_server->get_config()->clients_port_.load(std::memory_order_acquire))); +// boost::asio::connect(_client.next_layer(), _results); +// +// const auto _host = fmt::format("127.0.0.1:{}", std::to_string(_local_server->get_config()->clients_port_.load(std::memory_order_acquire))); +// _client.handshake(_host, "/"); +// +// boost::beast::flat_buffer _accepted_buffer; +// _client.read(_accepted_buffer); +// std::cout << boost::beast::make_printable(_accepted_buffer.data()) << std::endl; +// +// ASSERT_TRUE(_local_server->get_state()->get_clients().size() == 1); +// +// boost::system::error_code ec; +// _client.close(boost::beast::websocket::close_code::normal, ec); +// _client.next_layer().close(ec); +// +// // Wait for processing. +// std::this_thread::sleep_for(std::chrono::seconds(3)); +// +// ASSERT_TRUE(_local_server->get_state()->get_clients().size() == 0); +// } +// +// TEST_F(server_test, assert_local_server_can_handle_subscribe) { +// boost::asio::io_context _ioc; +// boost::asio::ip::tcp::resolver _resolver{make_strand(_ioc)}; +// boost::beast::websocket::stream _client{make_strand(_ioc)}; +// +// auto const _results = _resolver.resolve("127.0.0.1", std::to_string(_local_server->get_config()->clients_port_.load(std::memory_order_acquire))); +// boost::asio::connect(_client.next_layer(), _results); +// +// const auto _host = fmt::format("127.0.0.1:{}", std::to_string(_local_server->get_config()->clients_port_.load(std::memory_order_acquire))); +// _client.handshake(_host, "/"); +// +// boost::beast::flat_buffer _accepted_buffer; +// _client.read(_accepted_buffer); +// std::cout << boost::beast::make_printable(_accepted_buffer.data()) << std::endl; +// +// ASSERT_TRUE(_local_server->get_state()->get_clients().size() == 1); +// +// _client.write(boost::asio::buffer(std::string(serialize(boost::json::object{ +// {"transaction_id", to_string(boost::uuids::random_generator()())}, +// {"action", "subscribe"}, +// {"params", {{"channel", "welcome"}}}, +// })))); +// +// boost::beast::flat_buffer _buffer; +// +// _client.read(_buffer); +// +// std::cout << boost::beast::make_printable(_buffer.data()) << std::endl; +// +// boost::system::error_code ec; +// _client.close(boost::beast::websocket::close_code::normal, ec); +// _client.next_layer().close(ec); +// } +// +// TEST_F(server_test, assert_local_server_can_handle_publish) { +// boost::asio::io_context _ioc; +// boost::asio::ip::tcp::resolver _resolver{make_strand(_ioc)}; +// boost::beast::websocket::stream _local_client{make_strand(_ioc)}; +// boost::beast::websocket::stream _other_local_client{make_strand(_ioc)}; +// boost::beast::websocket::stream _remote_client{make_strand(_ioc)}; +// boost::beast::websocket::stream _other_remote_client{make_strand(_ioc)}; +// +// { +// auto const _results = _resolver.resolve("127.0.0.1", std::to_string(_local_server->get_config()->clients_port_.load(std::memory_order_acquire))); +// boost::asio::connect(_local_client.next_layer(), _results); +// boost::asio::connect(_other_local_client.next_layer(), _results); +// } +// +// { +// auto const _results = _resolver.resolve("127.0.0.1", std::to_string(_remote_server->get_config()->clients_port_.load(std::memory_order_acquire))); +// boost::asio::connect(_remote_client.next_layer(), _results); +// boost::asio::connect(_other_remote_client.next_layer(), _results); +// } +// +// const auto _host = fmt::format("127.0.0.1:{}", std::to_string(_local_server->get_config()->clients_port_.load(std::memory_order_acquire))); +// _local_client.handshake(_host, "/"); +// _other_local_client.handshake(_host, "/"); +// +// _remote_client.handshake(_host, "/"); +// _other_remote_client.handshake(_host, "/"); +// +// { +// boost::beast::flat_buffer _buffer; +// _local_client.read(_buffer); +// std::cout << "Local client should receive accepted ..." << std::endl; +// std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; +// } +// +// { +// boost::beast::flat_buffer _buffer; +// _other_local_client.read(_buffer); +// std::cout << "Other local client should receive accepted ..." << std::endl; +// std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; +// } +// +// { +// boost::beast::flat_buffer _buffer; +// _remote_client.read(_buffer); +// std::cout << "Remote client should receive accepted ..." << std::endl; +// std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; +// } +// +// { +// boost::beast::flat_buffer _buffer; +// _other_remote_client.read(_buffer); +// std::cout << "Other remote client should receive accepted ..." << std::endl; +// std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; +// _buffer.clear(); +// } +// +// _local_client.write(boost::asio::buffer(std::string(serialize(boost::json::object{ +// {"transaction_id", to_string(boost::uuids::random_generator()())}, +// {"action", "subscribe"}, +// {"params", {{"channel", "welcome"}}}, +// })))); +// +// { +// boost::beast::flat_buffer _buffer; +// _local_client.read(_buffer); +// std::cout << "Local client should receive subscribe ack ..." << std::endl; +// std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; +// _buffer.clear(); +// } +// +// _remote_client.write(boost::asio::buffer(std::string(serialize(boost::json::object{ +// {"transaction_id", to_string(boost::uuids::random_generator()())}, +// {"action", "subscribe"}, +// {"params", {{"channel", "welcome"}}}, +// })))); +// +// { +// boost::beast::flat_buffer _buffer; +// _remote_client.read(_buffer); +// std::cout << "Remote client should receive subscribe ack ..." << std::endl; +// std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; +// _buffer.clear(); +// } +// +// std::this_thread::sleep_for(std::chrono::seconds(1)); +// +// _other_local_client.write(boost::asio::buffer(std::string(serialize(boost::json::object{ +// {"transaction_id", to_string(boost::uuids::random_generator()())}, +// {"action", "publish"}, +// {"params", {{"channel", "welcome"}, {"payload", {{"message", "EHLO"}}}}}, +// })))); +// +// { +// boost::beast::flat_buffer _buffer; +// _other_local_client.read(_buffer); +// std::cout << "Other local client should receive publish ack ..." << std::endl; +// std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; +// _buffer.clear(); +// } +// +// std::this_thread::sleep_for(std::chrono::seconds(1)); +// +// { +// boost::beast::flat_buffer _buffer; +// _local_client.read(_buffer); +// std::cout << "Local client should receive publish message ..." << std::endl; +// std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; +// _buffer.clear(); +// } +// +// +// { +// boost::beast::flat_buffer _buffer; +// _remote_client.read(_buffer); +// std::cout << "Remote client should receive publish message ..." << std::endl; +// std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; +// } +// +// boost::system::error_code ec; +// _local_client.close(boost::beast::websocket::close_code::normal, ec); +// _other_local_client.close(boost::beast::websocket::close_code::normal, ec); +// _remote_client.close(boost::beast::websocket::close_code::normal, ec); +// _other_remote_client.close(boost::beast::websocket::close_code::normal, ec); +// } From 287bbd0635b7835fcd91dbf479869967477ec9ec Mon Sep 17 00:00:00 2001 From: Ian Torres Date: Sun, 4 Jan 2026 15:23:32 -0300 Subject: [PATCH 2/4] Testing client connectivity added. --- tests/server_test.cc | 73 ++++++++++++++++++++++++++++---------------- 1 file changed, 46 insertions(+), 27 deletions(-) diff --git a/tests/server_test.cc b/tests/server_test.cc index bd7e489..9d68d42 100644 --- a/tests/server_test.cc +++ b/tests/server_test.cc @@ -20,6 +20,7 @@ #include #include #include +#include TEST_F(server_test, assert_nodes_are_registered) { ASSERT_TRUE(server_b_->get_config()->registered_); @@ -29,33 +30,51 @@ TEST_F(server_test, assert_nodes_are_registered) { ASSERT_TRUE(server_c_->get_state()->get_sessions().size() == 2); } -// TEST_F(server_test, assert_local_server_accept_clients) { -// boost::asio::io_context _ioc; -// boost::asio::ip::tcp::resolver _resolver{make_strand(_ioc)}; -// boost::beast::websocket::stream _client{make_strand(_ioc)}; -// -// auto const _results = _resolver.resolve("127.0.0.1", std::to_string(_local_server->get_config()->clients_port_.load(std::memory_order_acquire))); -// boost::asio::connect(_client.next_layer(), _results); -// -// const auto _host = fmt::format("127.0.0.1:{}", std::to_string(_local_server->get_config()->clients_port_.load(std::memory_order_acquire))); -// _client.handshake(_host, "/"); -// -// boost::beast::flat_buffer _accepted_buffer; -// _client.read(_accepted_buffer); -// std::cout << boost::beast::make_printable(_accepted_buffer.data()) << std::endl; -// -// ASSERT_TRUE(_local_server->get_state()->get_clients().size() == 1); -// -// boost::system::error_code ec; -// _client.close(boost::beast::websocket::close_code::normal, ec); -// _client.next_layer().close(ec); -// -// // Wait for processing. -// std::this_thread::sleep_for(std::chrono::seconds(3)); -// -// ASSERT_TRUE(_local_server->get_state()->get_clients().size() == 0); -// } -// +TEST_F(server_test, assert_local_server_accept_clients) { + + for (auto &_server : { server_a_, server_b_, server_c_ }) { + boost::asio::io_context _ioc; + boost::asio::ip::tcp::resolver _resolver{make_strand(_ioc)}; + boost::beast::websocket::stream _client{make_strand(_ioc)}; + + auto const _results = _resolver.resolve("127.0.0.1", std::to_string(_server->get_config()->clients_port_.load(std::memory_order_acquire))); + boost::asio::connect(_client.next_layer(), _results); + + const auto _host = fmt::format("127.0.0.1:{}", std::to_string(_server->get_config()->clients_port_.load(std::memory_order_acquire))); + _client.handshake(_host, "/"); + + boost::beast::flat_buffer _accepted_buffer; + _client.read(_accepted_buffer); + auto _accepted_message = boost::beast::buffers_to_string(_accepted_buffer.data()); + + auto _accepted_object = boost::json::parse(_accepted_message); + + ASSERT_TRUE(_accepted_object.is_object()); + ASSERT_TRUE(_accepted_object.as_object().contains("action")); + ASSERT_TRUE(_accepted_object.as_object().at("action").is_string()); + ASSERT_EQ(_accepted_object.as_object().at("action").as_string(), "welcome"); + ASSERT_TRUE(_accepted_object.as_object().contains("transaction_id")); + ASSERT_TRUE(_accepted_object.as_object().at("transaction_id").is_string()); + ASSERT_TRUE(_accepted_object.as_object().contains("status")); + ASSERT_TRUE(_accepted_object.as_object().at("status").is_string()); + ASSERT_EQ(_accepted_object.as_object().at("status").as_string(), "success"); + ASSERT_TRUE(_accepted_object.as_object().contains("data")); + ASSERT_TRUE(_accepted_object.as_object().at("data").is_object()); + ASSERT_TRUE(_accepted_object.as_object().at("data").as_object().contains("client_id")); + ASSERT_TRUE(_accepted_object.as_object().at("data").as_object().at("client_id").is_string()); + + ASSERT_TRUE(_server->get_state()->get_clients().size() == 1); + + boost::system::error_code ec; + _client.close(boost::beast::websocket::close_code::normal, ec); + _client.next_layer().close(ec); + + std::this_thread::sleep_for(std::chrono::seconds(3)); + + ASSERT_TRUE(_server->get_state()->get_clients().size() == 0); + } +} + // TEST_F(server_test, assert_local_server_can_handle_subscribe) { // boost::asio::io_context _ioc; // boost::asio::ip::tcp::resolver _resolver{make_strand(_ioc)}; From 57ff62841dc3cf82d5692838ce910f943e54d0a9 Mon Sep 17 00:00:00 2001 From: Ian Torres Date: Sun, 4 Jan 2026 16:44:35 -0300 Subject: [PATCH 3/4] Fix. --- tests/server_test.cc | 346 +++++++++++++++++++++++-------------------- 1 file changed, 186 insertions(+), 160 deletions(-) diff --git a/tests/server_test.cc b/tests/server_test.cc index 9d68d42..3421cf1 100644 --- a/tests/server_test.cc +++ b/tests/server_test.cc @@ -22,7 +22,9 @@ #include #include -TEST_F(server_test, assert_nodes_are_registered) { +TEST_F(server_test, servers_are_registered) { + // Server isn't registered as isn't node + ASSERT_FALSE(server_a_->get_config()->registered_); ASSERT_TRUE(server_b_->get_config()->registered_); ASSERT_TRUE(server_c_->get_config()->registered_); ASSERT_TRUE(server_a_->get_state()->get_sessions().size() == 2); @@ -30,7 +32,7 @@ TEST_F(server_test, assert_nodes_are_registered) { ASSERT_TRUE(server_c_->get_state()->get_sessions().size() == 2); } -TEST_F(server_test, assert_local_server_accept_clients) { +TEST_F(server_test, servers_accept_clients) { for (auto &_server : { server_a_, server_b_, server_c_ }) { boost::asio::io_context _ioc; @@ -75,161 +77,185 @@ TEST_F(server_test, assert_local_server_accept_clients) { } } -// TEST_F(server_test, assert_local_server_can_handle_subscribe) { -// boost::asio::io_context _ioc; -// boost::asio::ip::tcp::resolver _resolver{make_strand(_ioc)}; -// boost::beast::websocket::stream _client{make_strand(_ioc)}; -// -// auto const _results = _resolver.resolve("127.0.0.1", std::to_string(_local_server->get_config()->clients_port_.load(std::memory_order_acquire))); -// boost::asio::connect(_client.next_layer(), _results); -// -// const auto _host = fmt::format("127.0.0.1:{}", std::to_string(_local_server->get_config()->clients_port_.load(std::memory_order_acquire))); -// _client.handshake(_host, "/"); -// -// boost::beast::flat_buffer _accepted_buffer; -// _client.read(_accepted_buffer); -// std::cout << boost::beast::make_printable(_accepted_buffer.data()) << std::endl; -// -// ASSERT_TRUE(_local_server->get_state()->get_clients().size() == 1); -// -// _client.write(boost::asio::buffer(std::string(serialize(boost::json::object{ -// {"transaction_id", to_string(boost::uuids::random_generator()())}, -// {"action", "subscribe"}, -// {"params", {{"channel", "welcome"}}}, -// })))); -// -// boost::beast::flat_buffer _buffer; -// -// _client.read(_buffer); -// -// std::cout << boost::beast::make_printable(_buffer.data()) << std::endl; -// -// boost::system::error_code ec; -// _client.close(boost::beast::websocket::close_code::normal, ec); -// _client.next_layer().close(ec); -// } -// -// TEST_F(server_test, assert_local_server_can_handle_publish) { -// boost::asio::io_context _ioc; -// boost::asio::ip::tcp::resolver _resolver{make_strand(_ioc)}; -// boost::beast::websocket::stream _local_client{make_strand(_ioc)}; -// boost::beast::websocket::stream _other_local_client{make_strand(_ioc)}; -// boost::beast::websocket::stream _remote_client{make_strand(_ioc)}; -// boost::beast::websocket::stream _other_remote_client{make_strand(_ioc)}; -// -// { -// auto const _results = _resolver.resolve("127.0.0.1", std::to_string(_local_server->get_config()->clients_port_.load(std::memory_order_acquire))); -// boost::asio::connect(_local_client.next_layer(), _results); -// boost::asio::connect(_other_local_client.next_layer(), _results); -// } -// -// { -// auto const _results = _resolver.resolve("127.0.0.1", std::to_string(_remote_server->get_config()->clients_port_.load(std::memory_order_acquire))); -// boost::asio::connect(_remote_client.next_layer(), _results); -// boost::asio::connect(_other_remote_client.next_layer(), _results); -// } -// -// const auto _host = fmt::format("127.0.0.1:{}", std::to_string(_local_server->get_config()->clients_port_.load(std::memory_order_acquire))); -// _local_client.handshake(_host, "/"); -// _other_local_client.handshake(_host, "/"); -// -// _remote_client.handshake(_host, "/"); -// _other_remote_client.handshake(_host, "/"); -// -// { -// boost::beast::flat_buffer _buffer; -// _local_client.read(_buffer); -// std::cout << "Local client should receive accepted ..." << std::endl; -// std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; -// } -// -// { -// boost::beast::flat_buffer _buffer; -// _other_local_client.read(_buffer); -// std::cout << "Other local client should receive accepted ..." << std::endl; -// std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; -// } -// -// { -// boost::beast::flat_buffer _buffer; -// _remote_client.read(_buffer); -// std::cout << "Remote client should receive accepted ..." << std::endl; -// std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; -// } -// -// { -// boost::beast::flat_buffer _buffer; -// _other_remote_client.read(_buffer); -// std::cout << "Other remote client should receive accepted ..." << std::endl; -// std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; -// _buffer.clear(); -// } -// -// _local_client.write(boost::asio::buffer(std::string(serialize(boost::json::object{ -// {"transaction_id", to_string(boost::uuids::random_generator()())}, -// {"action", "subscribe"}, -// {"params", {{"channel", "welcome"}}}, -// })))); -// -// { -// boost::beast::flat_buffer _buffer; -// _local_client.read(_buffer); -// std::cout << "Local client should receive subscribe ack ..." << std::endl; -// std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; -// _buffer.clear(); -// } -// -// _remote_client.write(boost::asio::buffer(std::string(serialize(boost::json::object{ -// {"transaction_id", to_string(boost::uuids::random_generator()())}, -// {"action", "subscribe"}, -// {"params", {{"channel", "welcome"}}}, -// })))); -// -// { -// boost::beast::flat_buffer _buffer; -// _remote_client.read(_buffer); -// std::cout << "Remote client should receive subscribe ack ..." << std::endl; -// std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; -// _buffer.clear(); -// } -// -// std::this_thread::sleep_for(std::chrono::seconds(1)); -// -// _other_local_client.write(boost::asio::buffer(std::string(serialize(boost::json::object{ -// {"transaction_id", to_string(boost::uuids::random_generator()())}, -// {"action", "publish"}, -// {"params", {{"channel", "welcome"}, {"payload", {{"message", "EHLO"}}}}}, -// })))); -// -// { -// boost::beast::flat_buffer _buffer; -// _other_local_client.read(_buffer); -// std::cout << "Other local client should receive publish ack ..." << std::endl; -// std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; -// _buffer.clear(); -// } -// -// std::this_thread::sleep_for(std::chrono::seconds(1)); -// -// { -// boost::beast::flat_buffer _buffer; -// _local_client.read(_buffer); -// std::cout << "Local client should receive publish message ..." << std::endl; -// std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; -// _buffer.clear(); -// } -// -// -// { -// boost::beast::flat_buffer _buffer; -// _remote_client.read(_buffer); -// std::cout << "Remote client should receive publish message ..." << std::endl; -// std::cout << boost::beast::make_printable(_buffer.data()) << std::endl << std::endl; -// } -// -// boost::system::error_code ec; -// _local_client.close(boost::beast::websocket::close_code::normal, ec); -// _other_local_client.close(boost::beast::websocket::close_code::normal, ec); -// _remote_client.close(boost::beast::websocket::close_code::normal, ec); -// _other_remote_client.close(boost::beast::websocket::close_code::normal, ec); -// } +TEST_F(server_test, server_can_handle_subscribe) { + boost::asio::io_context _ioc; + boost::asio::ip::tcp::resolver _resolver{make_strand(_ioc)}; + boost::beast::websocket::stream _client{make_strand(_ioc)}; + + auto const _results = _resolver.resolve("127.0.0.1", std::to_string(server_a_->get_config()->clients_port_.load(std::memory_order_acquire))); + boost::asio::connect(_client.next_layer(), _results); + + const auto _host = fmt::format("127.0.0.1:{}", std::to_string(server_a_->get_config()->clients_port_.load(std::memory_order_acquire))); + _client.handshake(_host, "/"); + + boost::beast::flat_buffer _accepted_buffer; + _client.read(_accepted_buffer); + ASSERT_TRUE(server_a_->get_state()->get_clients().size() == 1); + + _client.write(boost::asio::buffer(std::string(serialize(boost::json::object{ + {"transaction_id", to_string(boost::uuids::random_generator()())}, + {"action", "subscribe"}, + {"params", {{"channel", "welcome"}}}, + })))); + + boost::beast::flat_buffer _buffer; + _client.read(_buffer); + + + auto _subscribe_message = boost::beast::buffers_to_string(_accepted_buffer.data()); + auto _subscribe_object = boost::json::parse(_subscribe_message); + + ASSERT_TRUE(_subscribe_object.is_object()); + ASSERT_TRUE(_subscribe_object.as_object().contains("action")); + ASSERT_TRUE(_subscribe_object.as_object().at("action").is_string()); + ASSERT_EQ(_subscribe_object.as_object().at("action").as_string(), "welcome"); + ASSERT_TRUE(_subscribe_object.as_object().contains("transaction_id")); + ASSERT_TRUE(_subscribe_object.as_object().at("transaction_id").is_string()); + ASSERT_TRUE(_subscribe_object.as_object().contains("status")); + ASSERT_TRUE(_subscribe_object.as_object().at("status").is_string()); + ASSERT_EQ(_subscribe_object.as_object().at("status").as_string(), "success"); + ASSERT_TRUE(_subscribe_object.as_object().contains("data")); + ASSERT_TRUE(_subscribe_object.as_object().at("data").is_object()); + + boost::system::error_code ec; + _client.close(boost::beast::websocket::close_code::normal, ec); + _client.next_layer().close(ec); +} + +TEST_F(server_test, assert_local_server_can_handle_publish) { + boost::asio::io_context _ioc; + boost::asio::ip::tcp::resolver _resolver{make_strand(_ioc)}; + boost::beast::websocket::stream _client_a{make_strand(_ioc)}; + boost::beast::websocket::stream _client_b{make_strand(_ioc)}; + boost::beast::websocket::stream _client_c{make_strand(_ioc)}; + boost::beast::websocket::stream _client_d{make_strand(_ioc)}; + + { + auto const _results = _resolver.resolve("127.0.0.1", std::to_string(server_b_->get_config()->clients_port_.load(std::memory_order_acquire))); + boost::asio::connect(_client_a.next_layer(), _results); + boost::asio::connect(_client_b.next_layer(), _results); + } + + { + auto const _results = _resolver.resolve("127.0.0.1", std::to_string(server_c_->get_config()->clients_port_.load(std::memory_order_acquire))); + boost::asio::connect(_client_c.next_layer(), _results); + boost::asio::connect(_client_d.next_layer(), _results); + } + + const auto _server_b_host = fmt::format("127.0.0.1:{}", std::to_string(server_b_->get_config()->clients_port_.load(std::memory_order_acquire))); + _client_a.handshake(_server_b_host, "/"); + _client_b.handshake(_server_b_host, "/"); + + const auto _server_c_host = fmt::format("127.0.0.1:{}", std::to_string(server_c_->get_config()->clients_port_.load(std::memory_order_acquire))); + _client_c.handshake(_server_c_host, "/"); + _client_d.handshake(_server_c_host, "/"); + + for (const auto _client : { &_client_a, &_client_b, &_client_c, &_client_d }) { + boost::beast::flat_buffer _buffer; + _client->read(_buffer); + LOG_INFO("receiving client welcome ..."); + } + + _client_a.write(boost::asio::buffer(std::string(serialize(boost::json::object{ + {"transaction_id", to_string(boost::uuids::random_generator()())}, + {"action", "subscribe"}, + {"params", {{"channel", "welcome"}}}, + })))); + + { + boost::beast::flat_buffer _buffer; + _client_a.read(_buffer); + LOG_INFO("client A should receive subscribe ACK ... {}", boost::beast::buffers_to_string(_buffer.data())); + _buffer.clear(); + } + + _client_c.write(boost::asio::buffer(std::string(serialize(boost::json::object{ + {"transaction_id", to_string(boost::uuids::random_generator()())}, + {"action", "subscribe"}, + {"params", {{"channel", "welcome"}}}, + })))); + + { + boost::beast::flat_buffer _buffer; + _client_c.read(_buffer); + LOG_INFO("client C should receive subscribe ACK ... {}", boost::beast::buffers_to_string(_buffer.data())); + _buffer.clear(); + } + + std::this_thread::sleep_for(std::chrono::seconds(1)); + + _client_b.write(boost::asio::buffer(std::string(serialize(boost::json::object{ + {"transaction_id", to_string(boost::uuids::random_generator()())}, + {"action", "publish"}, + {"params", {{"channel", "welcome"}, {"payload", {{"message", "EHLO"}}}}}, + })))); + + { + boost::beast::flat_buffer _buffer; + _client_b.read(_buffer); + LOG_INFO("client B should receive publish ACK ... {}", boost::beast::buffers_to_string(_buffer.data())); + _buffer.clear(); + } + + std::this_thread::sleep_for(std::chrono::seconds(1)); + + { + boost::beast::flat_buffer _buffer; + _client_a.read(_buffer); + LOG_INFO("client A should receive publish MESSAGE ... {}", boost::beast::buffers_to_string(_buffer.data())); + + auto _publish_message = boost::beast::buffers_to_string(_buffer.data()); + auto _publish_object = boost::json::parse(_publish_message); + + ASSERT_TRUE(_publish_object.is_object()); + ASSERT_TRUE(_publish_object.as_object().contains("action")); + ASSERT_TRUE(_publish_object.as_object().at("action").is_string()); + ASSERT_EQ(_publish_object.as_object().at("action").as_string(), "publish"); + ASSERT_TRUE(_publish_object.as_object().contains("transaction_id")); + ASSERT_TRUE(_publish_object.as_object().at("transaction_id").is_string()); + ASSERT_TRUE(_publish_object.as_object().contains("params")); + ASSERT_TRUE(_publish_object.as_object().at("params").is_object()); + ASSERT_TRUE(_publish_object.as_object().at("params").as_object().contains("channel")); + ASSERT_TRUE(_publish_object.as_object().at("params").as_object().at("channel").is_string()); + ASSERT_EQ(_publish_object.as_object().at("params").as_object().at("channel").as_string(), "welcome"); + ASSERT_TRUE(_publish_object.as_object().at("params").as_object().contains("payload")); + ASSERT_TRUE(_publish_object.as_object().at("params").as_object().at("payload").is_object()); + ASSERT_TRUE(_publish_object.as_object().at("params").as_object().at("payload").as_object().contains("message")); + ASSERT_TRUE(_publish_object.as_object().at("params").as_object().at("payload").as_object().at("message").is_string()); + ASSERT_EQ(_publish_object.as_object().at("params").as_object().at("payload").as_object().at("message").as_string(), "EHLO"); + _buffer.clear(); + } + + + { + boost::beast::flat_buffer _buffer; + _client_c.read(_buffer); + LOG_INFO("client C should receive publish MESSAGE ... {}", boost::beast::buffers_to_string(_buffer.data())); + + auto _publish_message = boost::beast::buffers_to_string(_buffer.data()); + auto _publish_object = boost::json::parse(_publish_message); + + ASSERT_TRUE(_publish_object.is_object()); + ASSERT_TRUE(_publish_object.as_object().contains("action")); + ASSERT_TRUE(_publish_object.as_object().at("action").is_string()); + ASSERT_EQ(_publish_object.as_object().at("action").as_string(), "publish"); + ASSERT_TRUE(_publish_object.as_object().contains("transaction_id")); + ASSERT_TRUE(_publish_object.as_object().at("transaction_id").is_string()); + ASSERT_TRUE(_publish_object.as_object().contains("params")); + ASSERT_TRUE(_publish_object.as_object().at("params").is_object()); + ASSERT_TRUE(_publish_object.as_object().at("params").as_object().contains("channel")); + ASSERT_TRUE(_publish_object.as_object().at("params").as_object().at("channel").is_string()); + ASSERT_EQ(_publish_object.as_object().at("params").as_object().at("channel").as_string(), "welcome"); + ASSERT_TRUE(_publish_object.as_object().at("params").as_object().contains("payload")); + ASSERT_TRUE(_publish_object.as_object().at("params").as_object().at("payload").is_object()); + ASSERT_TRUE(_publish_object.as_object().at("params").as_object().at("payload").as_object().contains("message")); + ASSERT_TRUE(_publish_object.as_object().at("params").as_object().at("payload").as_object().at("message").is_string()); + ASSERT_EQ(_publish_object.as_object().at("params").as_object().at("payload").as_object().at("message").as_string(), "EHLO"); + } + + boost::system::error_code ec; + _client_a.close(boost::beast::websocket::close_code::normal, ec); + _client_b.close(boost::beast::websocket::close_code::normal, ec); + _client_c.close(boost::beast::websocket::close_code::normal, ec); + _client_d.close(boost::beast::websocket::close_code::normal, ec); +} From e8946cf2598e0540adf9ba4c9b58b675cc022e55 Mon Sep 17 00:00:00 2001 From: Ian Torres Date: Sun, 4 Jan 2026 17:11:29 -0300 Subject: [PATCH 4/4] Broadcast and send tested. --- tests/server_test.cc | 162 +++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 162 insertions(+) diff --git a/tests/server_test.cc b/tests/server_test.cc index 3421cf1..e87d143 100644 --- a/tests/server_test.cc +++ b/tests/server_test.cc @@ -259,3 +259,165 @@ TEST_F(server_test, assert_local_server_can_handle_publish) { _client_c.close(boost::beast::websocket::close_code::normal, ec); _client_d.close(boost::beast::websocket::close_code::normal, ec); } + +TEST_F(server_test, assert_local_server_can_handle_broadcast) { + boost::asio::io_context _ioc; + boost::asio::ip::tcp::resolver _resolver{make_strand(_ioc)}; + boost::beast::websocket::stream _client_a{make_strand(_ioc)}; + boost::beast::websocket::stream _client_b{make_strand(_ioc)}; + boost::beast::websocket::stream _client_c{make_strand(_ioc)}; + boost::beast::websocket::stream _client_d{make_strand(_ioc)}; + + { + auto const _results = _resolver.resolve("127.0.0.1", std::to_string(server_b_->get_config()->clients_port_.load(std::memory_order_acquire))); + boost::asio::connect(_client_a.next_layer(), _results); + boost::asio::connect(_client_b.next_layer(), _results); + } + + { + auto const _results = _resolver.resolve("127.0.0.1", std::to_string(server_c_->get_config()->clients_port_.load(std::memory_order_acquire))); + boost::asio::connect(_client_c.next_layer(), _results); + boost::asio::connect(_client_d.next_layer(), _results); + } + + const auto _server_b_host = fmt::format("127.0.0.1:{}", std::to_string(server_b_->get_config()->clients_port_.load(std::memory_order_acquire))); + _client_a.handshake(_server_b_host, "/"); + _client_b.handshake(_server_b_host, "/"); + + const auto _server_c_host = fmt::format("127.0.0.1:{}", std::to_string(server_c_->get_config()->clients_port_.load(std::memory_order_acquire))); + _client_c.handshake(_server_c_host, "/"); + _client_d.handshake(_server_c_host, "/"); + + for (const auto _client : { &_client_a, &_client_b, &_client_c, &_client_d }) { + boost::beast::flat_buffer _buffer; + _client->read(_buffer); + LOG_INFO("receiving client welcome ..."); + } + + _client_a.write(boost::asio::buffer(std::string(serialize(boost::json::object{ + {"transaction_id", to_string(boost::uuids::random_generator()())}, + {"action", "broadcast"}, + {"params", {{"payload", {{"message", "EHLO"}}}}}, + })))); + + { + boost::beast::flat_buffer _buffer; + _client_a.read(_buffer); + LOG_INFO("client A should receive broadcast ACK ... {}", boost::beast::buffers_to_string(_buffer.data())); + _buffer.clear(); + } + + std::this_thread::sleep_for(std::chrono::seconds(1)); + + for (const auto _client : { &_client_b, &_client_c, &_client_d }) + { + boost::beast::flat_buffer _buffer; + _client->read(_buffer); + LOG_INFO("client should receive broadcast MESSAGE ... {}", boost::beast::buffers_to_string(_buffer.data())); + + auto _publish_message = boost::beast::buffers_to_string(_buffer.data()); + auto _publish_object = boost::json::parse(_publish_message); + + ASSERT_TRUE(_publish_object.is_object()); + ASSERT_TRUE(_publish_object.as_object().contains("action")); + ASSERT_TRUE(_publish_object.as_object().at("action").is_string()); + ASSERT_EQ(_publish_object.as_object().at("action").as_string(), "broadcast"); + ASSERT_TRUE(_publish_object.as_object().contains("transaction_id")); + ASSERT_TRUE(_publish_object.as_object().at("transaction_id").is_string()); + ASSERT_TRUE(_publish_object.as_object().contains("params")); + ASSERT_TRUE(_publish_object.as_object().at("params").is_object()); + ASSERT_TRUE(_publish_object.as_object().at("params").as_object().contains("payload")); + ASSERT_TRUE(_publish_object.as_object().at("params").as_object().at("payload").is_object()); + ASSERT_TRUE(_publish_object.as_object().at("params").as_object().at("payload").as_object().contains("message")); + ASSERT_TRUE(_publish_object.as_object().at("params").as_object().at("payload").as_object().at("message").is_string()); + ASSERT_EQ(_publish_object.as_object().at("params").as_object().at("payload").as_object().at("message").as_string(), "EHLO"); + _buffer.clear(); + } + + boost::system::error_code ec; + _client_a.close(boost::beast::websocket::close_code::normal, ec); + _client_b.close(boost::beast::websocket::close_code::normal, ec); + _client_c.close(boost::beast::websocket::close_code::normal, ec); + _client_d.close(boost::beast::websocket::close_code::normal, ec); +} + + +TEST_F(server_test, assert_local_server_can_handle_send) { + boost::asio::io_context _ioc; + boost::asio::ip::tcp::resolver _resolver{make_strand(_ioc)}; + boost::beast::websocket::stream _client_a{make_strand(_ioc)}; + boost::beast::websocket::stream _client_b{make_strand(_ioc)}; + + { + auto const _results = _resolver.resolve("127.0.0.1", std::to_string(server_b_->get_config()->clients_port_.load(std::memory_order_acquire))); + boost::asio::connect(_client_a.next_layer(), _results); + } + + { + auto const _results = _resolver.resolve("127.0.0.1", std::to_string(server_c_->get_config()->clients_port_.load(std::memory_order_acquire))); + boost::asio::connect(_client_b.next_layer(), _results); + } + + const auto _server_b_host = fmt::format("127.0.0.1:{}", std::to_string(server_b_->get_config()->clients_port_.load(std::memory_order_acquire))); + _client_a.handshake(_server_b_host, "/"); + + const auto _server_c_host = fmt::format("127.0.0.1:{}", std::to_string(server_c_->get_config()->clients_port_.load(std::memory_order_acquire))); + _client_b.handshake(_server_c_host, "/"); + + boost::beast::flat_buffer _client_a_accepted_buffer; + _client_a.read(_client_a_accepted_buffer); + auto _accepted_a_message = boost::beast::buffers_to_string(_client_a_accepted_buffer.data()); + auto _accepted_a_object = boost::json::parse(_accepted_a_message); + + auto _client_a_id = _accepted_a_object.as_object().at("data").at("client_id").as_string(); + + boost::beast::flat_buffer _client_b_accepted_buffer; + _client_b.read(_client_b_accepted_buffer); + auto _accepted_b_message = boost::beast::buffers_to_string(_client_b_accepted_buffer.data()); + auto _accepted_b_object = boost::json::parse(_accepted_b_message); + + auto _client_b_id = _accepted_b_object.as_object().at("data").at("client_id").as_string(); + + _client_a.write(boost::asio::buffer(std::string(serialize(boost::json::object{ + {"transaction_id", to_string(boost::uuids::random_generator()())}, + {"action", "send"}, + {"params", { + {"to_client_id", _client_b_id}, + {"payload", {{"message", "EHLO"}}}}}, + })))); + + { + boost::beast::flat_buffer _buffer; + _client_a.read(_buffer); + LOG_INFO("client A should receive send ACK ... {}", boost::beast::buffers_to_string(_buffer.data())); + _buffer.clear(); + } + + std::this_thread::sleep_for(std::chrono::seconds(1)); + + boost::beast::flat_buffer _buffer; + _client_b.read(_buffer); + LOG_INFO("client should receive send MESSAGE ... {}", boost::beast::buffers_to_string(_buffer.data())); + + auto _send_message = boost::beast::buffers_to_string(_buffer.data()); + auto _send_object = boost::json::parse(_send_message); + + ASSERT_TRUE(_send_object.is_object()); + ASSERT_TRUE(_send_object.as_object().contains("action")); + ASSERT_TRUE(_send_object.as_object().at("action").is_string()); + ASSERT_EQ(_send_object.as_object().at("action").as_string(), "send"); + ASSERT_TRUE(_send_object.as_object().contains("transaction_id")); + ASSERT_TRUE(_send_object.as_object().at("transaction_id").is_string()); + ASSERT_TRUE(_send_object.as_object().contains("params")); + ASSERT_TRUE(_send_object.as_object().at("params").is_object()); + ASSERT_TRUE(_send_object.as_object().at("params").as_object().contains("payload")); + ASSERT_TRUE(_send_object.as_object().at("params").as_object().at("payload").is_object()); + ASSERT_TRUE(_send_object.as_object().at("params").as_object().at("payload").as_object().contains("message")); + ASSERT_TRUE(_send_object.as_object().at("params").as_object().at("payload").as_object().at("message").is_string()); + ASSERT_EQ(_send_object.as_object().at("params").as_object().at("payload").as_object().at("message").as_string(), "EHLO"); + + + boost::system::error_code ec; + _client_a.close(boost::beast::websocket::close_code::normal, ec); + _client_b.close(boost::beast::websocket::close_code::normal, ec); +}