You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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 &amp;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
相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.08 10:57:38