Merge pull request #167 from qicosmos/check_port

This commit is contained in:
qicosmos
2025-09-08 17:54:12 +08:00
committed by GitHub
2 changed files with 96 additions and 7 deletions
+71 -7
View File
@@ -17,11 +17,10 @@ class rpc_server : private asio::noncopyable {
public:
rpc_server(unsigned short port, size_t size, size_t timeout_seconds = 15,
size_t check_seconds = 10)
: io_service_pool_(size), acceptor_(io_service_pool_.get_io_service(),
tcp::endpoint(tcp::v4(), port)),
: io_service_pool_(size), acceptor_(io_service_pool_.get_io_service()),
timeout_seconds_(timeout_seconds), check_seconds_(check_seconds),
signals_(io_service_pool_.get_io_service()) {
do_accept();
signals_(io_service_pool_.get_io_service()),
port_(std::to_string(port)) {
check_thread_ = std::make_shared<std::thread>([this] { clean(); });
pub_sub_thread_ =
std::make_shared<std::thread>([this] { clean_sub_pub(); });
@@ -33,6 +32,12 @@ public:
do_await_stop();
}
rpc_server(std::string address, unsigned short port, size_t size,
size_t timeout_seconds = 15, size_t check_seconds = 10)
: rpc_server(port, size, timeout_seconds, check_seconds) {
address_ = std::move(address);
}
rpc_server(unsigned short port, size_t size, ssl_configure ssl_conf,
size_t timeout_seconds = 15, size_t check_seconds = 10)
: rpc_server(port, size, timeout_seconds, check_seconds) {
@@ -46,11 +51,24 @@ public:
~rpc_server() { stop(); }
void async_run() {
thd_ = std::make_shared<std::thread>([this] { io_service_pool_.run(); });
std::error_code async_run() {
auto ec = listen();
if (!ec) {
do_accept();
thd_ = std::make_shared<std::thread>([this] { io_service_pool_.run(); });
}
return ec;
}
void run() { io_service_pool_.run(); }
std::error_code run() {
auto ec = listen();
if (!ec) {
do_accept();
io_service_pool_.run();
}
return ec;
}
template <bool is_pub = false, typename Function>
void register_handler(std::string const &name, const Function &f) {
@@ -132,6 +150,50 @@ private:
});
}
std::error_code listen() {
using asio::ip::tcp;
asio::error_code ec;
asio::ip::tcp::resolver resolver(acceptor_.get_executor());
auto endpoints = resolver.resolve(address_, port_, ec);
if (ec) {
return ec;
}
auto it = endpoints.begin();
if (it == endpoints.end()) {
return std::make_error_code(std::errc::bad_address);
}
auto endpoint = it->endpoint();
acceptor_.open(endpoint.protocol(), ec);
if (ec) {
return ec;
}
#ifdef __GNUC__
acceptor_.set_option(tcp::acceptor::reuse_address(true), ec);
#endif
acceptor_.bind(endpoint, ec);
if (ec) {
std::error_code ignore;
acceptor_.cancel(ignore);
acceptor_.close(ignore);
return ec;
}
#ifdef _MSC_VER
acceptor_.set_option(tcp::acceptor::reuse_address(true));
#endif
acceptor_.listen(asio::socket_base::max_listen_connections, ec);
if (ec) {
std::error_code ignore;
acceptor_.cancel(ignore);
acceptor_.close(ignore);
return ec;
}
return ec;
}
void clean() {
while (!stop_check_) {
std::unique_lock<std::mutex> lock(mtx_);
@@ -259,6 +321,8 @@ private:
std::shared_ptr<connection> conn_;
std::shared_ptr<std::thread> thd_;
std::size_t timeout_seconds_;
std::string address_ = "0.0.0.0";
std::string port_ = "";
std::unordered_map<int64_t, std::shared_ptr<connection>> connections_;
int64_t conn_id_ = 0;
+25
View File
@@ -109,6 +109,31 @@ TEST_CASE("test_client_default_constructor") {
CHECK_EQ(result, 3);
}
TEST_CASE("test start some servers with same port") {
rpc_server server1(9000, 1);
rpc_server server2(9000, 1);
auto ec1 = server1.async_run();
CHECK(ec1 == std::error_code{});
auto ec2 = server2.async_run();
CHECK(ec2);
std::cout << ec2.message() << "\n";
rpc_server server3(9000, 1);
auto ec3 = server3.async_run();
CHECK(ec3);
std::cout << ec3.message() << "\n";
}
TEST_CASE("test start server with local ip") {
rpc_server server1("0.0.0.0", 9000, 1);
auto ec1 = server1.async_run();
CHECK(ec1 == std::error_code{});
rpc_server server2("11.11.11.11", 9000, 1);
auto ec2 = server2.async_run();
CHECK(ec2);
std::cout << ec2.message() << "\n"; // address not available
}
TEST_CASE("test_constructor_with_language") {
rpc_server server(9000, std::thread::hardware_concurrency());
dummy d;