mirror of
https://github.com/qicosmos/rest_rpc.git
synced 2026-08-29 08:34:47 +08:00
set address; set reuse port
This commit is contained in:
@@ -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;
|
||||
|
||||
@@ -109,6 +109,29 @@ 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::errc::address_in_use);
|
||||
rpc_server server3(9000, 1);
|
||||
auto ec3 = server3.async_run();
|
||||
CHECK(ec3 == std::errc::address_in_use);
|
||||
}
|
||||
|
||||
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;
|
||||
|
||||
Reference in New Issue
Block a user