Boost.Asio异步客户端与服务器代理实现:跨实例套接字访问问题
问题描述
我正在实现一个客户端与服务器的代理系统,该代理负责接收客户端命令并转发至服务器,反之亦然。但在跨实例访问客户端套接字以向客户端转发消息(以及反向操作)时遇到了挑战,目前正在学习Boost.Asio的特性,恳请提供在服务器与客户端实例间共享套接字的最有效方案。当前实现中,cmd_handler::read_cmd_done接口内client_socket_.is_open()返回0而非1。
当前实现代码
头文件 proxy.h
#pragma once #ifndef __OAMIP_PROXY_H__ #define __OAMIP_PROXY_H__ #include <iostream> #include <boost/asio.hpp> #include <thread> #include <vector> #include <functional> #include <deque> #include "config.h" namespace asio = boost::asio; class cmd_handler: public std::enable_shared_from_this<cmd_handler> { public: cmd_handler(asio::io_context &io_context, AppConfig* appConfig); ~cmd_handler(); asio::ip::tcp::socket &socket(); asio::ip::tcp::socket &c_socket(); asio::ip::tcp::endpoint remote_endpoint(); asio::ip::tcp::endpoint c_remote_endpoint(); asio::io_context &m_io_context(); void start(); void read_cmd(); void read_cmd_done(boost::system::error_code const &ec, std::size_t bytes_transferred); void read_cmd_client(); void read_cmd_client_done(boost::system::error_code const &ec, std::size_t bytes_transferred); private: asio::io_context& io_context_; asio::ip::tcp::socket server_socket_; asio::ip::tcp::socket client_socket_; asio::io_context::strand write_strand_; asio::streambuf in_packet_; std::deque<std::string> send_cmd_queue; std::mutex queue_mutex_; // Added for thread safety AppConfig *app_config; }; class ProxyServer { using shared_handler_t = std::shared_ptr<cmd_handler>; public: ProxyServer(int thread_count, AppConfig* appConfig);; ~ProxyServer(); void start_server(std::string ip_addr, int port); void start_client(std::string ip_addr, int port); void handle_new_connection(shared_handler_t handler, boost::system::error_code const &ec); private: asio::io_context io_context_; int thread_count_; asio::ip::tcp::acceptor acceptor_; asio::ip::tcp::resolver resolver_; std::vector<std::thread> thread_pool_; asio::streambuf buffer_; AppConfig *app_config; }; #endif // __OAMIP_PROXY_H__
实现文件 proxy.cpp
#include "proxy.h" #include "parser.cpp" #include "send.cpp" ProxyServer::ProxyServer(int thread_count, AppConfig *appconfig) : thread_count_(thread_count), acceptor_(io_context_), resolver_(io_context_), thread_pool_(), app_config(appconfig) { } /*******************************************************************************************/ ProxyServer::~ProxyServer() { // Stop and join the io_context to prevent memory leaks io_context_.stop(); for (auto &thread : thread_pool_) { thread.join(); } } /*******************************************************************************************/ void ProxyServer::start_server(std::string ip_addr, int port) { std::cout << "Starting Server, IP:" << ip_addr << ", Port:" << port << std::endl; auto handler = std::make_shared<cmd_handler>(io_context_, app_config); asio::ip::address_v4 ipv4_address = asio::ip::address_v4::from_string(ip_addr); asio::ip::tcp::endpoint endpoint(ipv4_address, port); acceptor_.open(endpoint.protocol()); acceptor_.set_option(asio::ip::tcp::acceptor::reuse_address(true)); acceptor_.bind(endpoint); acceptor_.listen(); std::cout << "Start listening" << std::endl; try { acceptor_.async_accept(handler->socket(), [=](auto ec) { if(ec) std::cerr << "Error accepting connection from : " << ec.message() << std::endl; else handle_new_connection(handler, ec); }); } catch (const std::exception &e) { std::cerr << "Exception caught: " << e.what() << std::endl; } io_context_.run(); // start pool of threads to process the asio events /* for (int i = 0; i < thread_count_; ++i) { thread_pool_.emplace_back([=] { io_context_.run(); }); } for (auto &thread : thread_pool_) { thread.join(); } */ } void ProxyServer::start_client(std::string ip_addr, int port) { std::cout << "Starting client" << std::endl; auto handler = std::make_shared<cmd_handler>(io_context_, app_config); asio::ip::tcp::resolver::query query(ip_addr, std::to_string(port)); asio::ip::tcp::resolver::iterator endpoint_iterator = resolver_.resolve(query); asio::connect(handler->c_socket(), endpoint_iterator); std::cout << "Connected to server on port " << port << std::endl; handler->read_cmd_client(); io_context_.run(); } /*******************************************************************************************/ void ProxyServer::handle_new_connection(shared_handler_t handler, boost::system::error_code const &ec) { std::cout << "Handle connection" << std::endl; if (ec) { std::cerr << "Error accepting connection from client: " << ec.message() << std::endl; return; } handler->read_cmd(); auto new_handler = std::make_shared<cmd_handler>(io_context_, app_config); acceptor_.async_accept(new_handler->socket(), [=](auto ec) { handle_new_connection(new_handler, ec); }); } /*******************************************************************************************/ cmd_handler::cmd_handler(asio::io_context &io_context, AppConfig *appConfig) : io_context_(io_context), server_socket_(io_context), client_socket_(io_context),write_strand_(io_context), app_config(appConfig) { } /*******************************************************************************************/ cmd_handler::~cmd_handler() { // Explicitly clear the buffer to release the allocated memory in_packet_.consume(in_packet_.size()); } /*******************************************************************************************/ asio::ip::tcp::socket &cmd_handler::socket() { return server_socket_; } asio::ip::tcp::socket &cmd_handler::c_socket() { return client_socket_; } asio::ip::tcp::endpoint cmd_handler::remote_endpoint() { return server_socket_.remote_endpoint(); } asio::ip::tcp::endpoint cmd_handler::c_remote_endpoint() { return client_socket_.remote_endpoint(); } asio::io_context &cmd_handler::m_io_context() { return io_context_; } /*******************************************************************************************/ void cmd_handler::start() { read_cmd(); } /*******************************************************************************************/ void cmd_handler::read_cmd() { auto remote_ip = remote_endpoint().address().to_string(); std::cout << "Read command from " << remote_ip << std::endl; asio::async_read_until(server_socket_, in_packet_, '\n', [me = shared_from_this()](boost::system::error_code const &ec, std::size_t bytes_xfer) { if (ec == asio::error::eof) { std::cout << "Connection closed by client:" << std::endl; return; // No need to read further; the connection is closed. } else if (ec) { std::cerr << "Error in async_read_until: " << ec.message() << std::endl; return; } else { me->read_cmd_done(ec, bytes_xfer); } }); } void cmd_handler::read_cmd_client() { std::cout << "Read cmd client" << std::endl; auto remote_ip = c_remote_endpoint().address().to_string(); std::cout << client_socket_.is_open() << std::endl; std::cout << "Read command from " << remote_ip << std::endl; asio::async_read_until(client_socket_, in_packet_, '\n', [me = shared_from_this()](boost::system::error_code const &ec, std::size_t bytes_xfer) { if (ec == asio::error::eof) { std::cout << "Connection closed by client:" << std::endl; return; // No need to read further; the connection is closed. } else if (ec) { std::cerr << "Error in async_read_until: " << ec.message() << std::endl; return; } else { me->read_cmd_client_done(ec, bytes_xfer); } }); } void cmd_handler::read_cmd_client_done(boost::system::error_code const &ec, std::size_t bytes_transferred) { auto remote_ip = remote_endpoint().address().to_string(); if (ec == asio::error::eof) { std::cout << "Connection closed by client:" << remote_ip << std::endl; return; // No need to read further; the connection is closed. } else if (ec) { std::cerr << "Error accepting packet from the client: " << remote_ip << "," << ec.message() << std::endl; return; } std::string command(buffers_begin(in_packet_.data()), buffers_begin(in_packet_.data()) + bytes_transferred); in_packet_.consume(bytes_transferred); std::cout << "Connected server IP: " << remote_ip << std::endl; std::cout << "command:" << command << std::endl; Parser parser(app_config); std::string recv_cmd = parser.process_cmd(command); //CmdSender sender(server_socket_); //sender.send_cmd(recv_cmd); read_cmd(); }; /*******************************************************************************************/ void cmd_handler::read_cmd_done(boost::system::error_code const &ec, std::size_t bytes_transferred) { auto remote_ip = remote_endpoint().address().to_string(); if (ec == asio::error::eof) { std::cout << "Connection closed by client:" << remote_ip << std::endl; return; // No need to read further; the connection is closed. } else if (ec) { std::cerr << "Error accepting packet from the client: " << remote_ip << "," << ec.message() << std::endl; return; } std::string command(buffers_begin(in_packet_.data()), buffers_begin(in_packet_.data()) + bytes_transferred); in_packet_.consume(bytes_transferred); std::cout << "Connected client IP: " << remote_ip << std::endl; std::cout << "command:" << command << std::endl; Parser parser(app_config); std::string recv_cmd = parser.process_cmd(command); std::cout << client_socket_.is_open() <<c_socket().is_open() << socket().is_open()<< std::endl; CmdSender sender(client_socket_); sender.send_cmd(recv_cmd); read_cmd(); }; /*******************************************************************************************/ int main() { try { int thread_count = 1; int port = appConfig.serverConfig.port; int s_port = appConfig.clientConfig.port; std::vector<std::thread> thread_pool; std::string ip_addr = appConfig.serverConfig.ip; // start pool of threads to process the asio events for (int i = 0; i < thread_count; ++i) { thread_pool.emplace_back([&]() { proxy_server.start_server(ip_addr, port); }); } for (int i = 0; i < thread_count; ++i) { thread_pool.emplace_back([&]() { proxy_server.start_client(ip_addr, s_port);}); } for (auto &thread : thread_pool) { thread.join(); } } catch (const std::exception &e) { std::cerr << "Exception caught: " << e.what() << std::endl; } return 0; }
解决方案
核心问题分析
你的代码中存在两个完全独立的cmd_handler实例:一个在start_server中创建,用于接收客户端连接;另一个在start_client中创建,用于连接目标服务器。这两个实例的套接字没有任何关联,所以当服务器端的cmd_handler访问自身的client_socket_时,该套接字从未被初始化连接,自然返回is_open() == false。
正确的代理模型实现
代理的核心逻辑是:每接受一个客户端连接,就创建一个对应的到目标服务器的连接,将两个套接字绑定到同一个cmd_handler实例中,实现双向数据转发。
1. 重构连接处理逻辑
删除ProxyServer::start_client方法,改为在handle_new_connection中,当接受客户端连接后,主动发起对目标服务器的异步连接:
void ProxyServer::handle_new_connection(shared_handler_t handler, boost::system::error_code const &ec) { std::cout << "Handle connection" << std::endl; if (ec) { std::cerr << "Error accepting connection from client: " << ec.message() << std::endl; return; } // 解析目标服务器地址并发起异步连接 asio::ip::tcp::resolver::query query(app_config->clientConfig.ip, std::to_string(app_config->clientConfig.port)); resolver_.async_resolve(query, [handler, this](boost::system::error_code ec, asio::ip::tcp::resolver::iterator endpoints) { if (ec) { std::cerr << "Resolve target server failed: " << ec.message() << std::endl; handler->socket().close(); return; } // 异步连接目标服务器 asio::async_connect(handler->c_socket(), endpoints, [handler](boost::system::error_code ec) { if (ec) { std::cerr << "Connect to target server failed: " << ec.message() << std::endl; handler->socket().close(); return; } // 启动双向读写 handler->read_cmd(); // 读取客户端数据 handler->read_cmd_client();// 读取目标服务器数据 }); }); // 继续监听新的客户端连接 auto new_handler = std::make_shared<cmd_handler>(io_context_, app_config); acceptor_.async_accept(new_handler->socket(), [=](auto ec) { handle_new_connection(new_handler, ec); }); }
2. 修复数据转发逻辑
添加异步写入方法,使用strand保证线程安全(避免多线程并发写入套接字):
// 在cmd_handler类中添加以下方法声明 void write_to_server(const std::string& data); void do_write_server(); void write_to_client(const std::string& data); void do_write_client(); // 实现异步写入逻辑 void cmd_handler::write_to_server(const std::string& data) { asio::post(write_strand_, [this, data]() { bool write_in_progress = !send_cmd_queue.empty(); send_cmd_queue.push_back(data); if (!write_in_progress) { do_write_server(); } }); } void cmd_handler::do_write_server() { asio::async_write(server_socket_, asio::buffer(send_cmd_queue.front()), asio::bind_executor(write_strand_, [this](boost::system::error_code ec, std::size_t) { if (!ec) { send_cmd_queue.pop_front(); if (!send_cmd_queue.empty()) { do_write_server(); } } else { std::cerr << "Write to server failed: " << ec.message() << std::endl; server_socket_.close(); client_socket_.close(); } })); } void cmd_handler::write_to_client(const std::string& data) { asio::post(write_strand_, [this, data]() { bool write_in_progress = !send_cmd_queue.empty(); send_cmd_queue.push_back(data); if (!write_in_progress) { do_write_client(); } }); } void cmd_handler::do_write_client() { asio::async_write(client_socket_, asio::buffer(s
相关产品推荐
相关产品推荐

