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

Boost SSL WebSocket服务器同时读写及异常捕获问题求助

Boost 1.72异步SSL WebSocket服务器问题修复

问题背景

基于Boost 1.72官方异步SSL WebSocket服务器示例开发时,遇到以下问题:

  • 在on_read回调中通过循环调用ws_.write多次发送数据时,无法接收同一客户端的后续消息
  • 尝试改用async_read时程序崩溃
  • 需要捕获boost::wrapexcept<boost::system::system_error>、stream truncated、Broken pipe等异常,避免程序终止

问题原因分析

  1. 同步操作阻塞IO线程:回调中使用同步ws_.write和sleep(2)会直接阻塞当前IO线程,导致io_context无法调度后续的异步async_read操作,因此无法接收客户端消息。
  2. 异步操作调用方式错误:直接在回调中循环发起异步操作如果没有正确管理状态,容易导致对象生命周期问题或重复操作,引发崩溃。
  3. 异常未正确捕获处理:官方示例的错误处理仅处理error_code,未捕获抛出的异常,导致程序终止。

修改后的完整代码

基于官方示例,修改session类的相关逻辑,实现异步链式发送、非阻塞IO、异常捕获:

#include <boost/beast.hpp>
#include <boost/asio.hpp>
#include <boost/asio/ssl.hpp>
#include <iostream>
#include <memory>
#include <string>
#include <chrono>

namespace beast = boost::beast;
namespace http = beast::http;
namespace websocket = beast::websocket;
namespace net = boost::asio;
namespace ssl = net::ssl;
using tcp = boost::asio::ip::tcp;

// 会话类,管理单个WebSocket连接
class session : public std::enable_shared_from_this<session>
{
    websocket::stream<ssl::stream<tcp::socket>> ws_;
    beast::flat_buffer buffer_;
    const std::string send_msg_ = "{\"dl_rate\":1029.360857,\"ul_rate\":7426.864}";
    short send_count_ = 20; // 需要发送的次数

public:
    explicit session(tcp::socket socket, ssl::context& ctx)
        : ws_(std::move(socket), ctx)
    {
    }

    // 启动会话
    void run()
    {
        // 设置SSL握手回调
        ws_.next_layer().async_handshake(ssl::stream_base::server,
            beast::bind_front_handler(
                &session::on_handshake,
                shared_from_this()));
    }

    void on_handshake(beast::error_code ec)
    {
        if(ec)
            return fail(ec, "handshake");

        // 接受WebSocket握手
        ws_.async_accept(
            beast::bind_front_handler(
                &session::on_accept,
                shared_from_this()));
    }

    void on_accept(beast::error_code ec)
    {
        if(ec)
            return fail(ec, "accept");

        // 开始读取消息
        do_read();
    }

    void do_read()
    {
        // 异步读取消息,捕获可能的异常
        ws_.async_read(
            buffer_,
            beast::bind_front_handler(
                &session::on_read,
                shared_from_this()));
    }

    void on_read(beast::error_code ec, std::size_t bytes_transferred)
    {
        boost::ignore_unused(bytes_transferred);

        // 处理关闭错误
        if(ec == websocket::error::closed)
            return;

        if(ec)
            return fail(ec, "read");

        // 打印收到的消息
        std::cout << "收到客户端消息: " << beast::make_printable(buffer_.data()) << std::endl;
        buffer_.consume(buffer_.size()); // 清空缓冲区

        // 设置文本模式
        ws_.text(ws_.got_text());

        // 重置发送计数器,开始异步发送链
        send_count_ = 20;
        do_write();
    }

    void do_write()
    {
        if(send_count_ <= 0)
        {
            // 发送完成,继续读取客户端消息
            do_read();
            return;
        }

        // 异步发送消息,捕获异常
        auto self(shared_from_this());
        ws_.async_write(
            net::buffer(send_msg_),
            [self](beast::error_code ec, std::size_t bytes_transferred)
            {
                boost::ignore_unused(bytes_transferred);

                if(ec)
                {
                    self->fail(ec, "write");
                    return;
                }

                // 计数器减1,延迟2秒后发送下一条(用异步定时器,避免阻塞)
                --self->send_count_;
                net::steady_timer timer(self->ws_.get_executor(), std::chrono::seconds(2));
                timer.async_wait(
                    [self](beast::error_code ec)
                    {
                        if(!ec)
                            self->do_write();
                    });
            });
    }

    // 错误处理函数,捕获异常并避免程序终止
    void fail(beast::error_code ec, char const* what)
    {
        // 处理特定错误码
        if(ec == net::error::broken_pipe || ec.message() == "stream truncated")
        {
            std::cerr << "连接异常: " << what << ": " << ec.message() << std::endl;
        }
        else if(ec)
        {
            std::cerr << "错误: " << what << ": " << ec.message() << std::endl;
        }

        // 捕获可能抛出的异常
        try
        {
            // 尝试正常关闭连接
            ws_.close(websocket::close_code::normal);
        }
        catch(const boost::wrapexcept<boost::system::system_error>& e)
        {
            std::cerr << "捕获异常: " << e.what() << std::endl;
        }
        catch(const std::exception& e)
        {
            std::cerr << "捕获未知异常: " << e.what() << std::endl;
        }
    }
};

// 监听类,接受新连接
class listener : public std::enable_shared_from_this<listener>
{
    net::io_context& ioc_;
    ssl::context& ctx_;
    tcp::acceptor acceptor_;

public:
    listener(net::io_context& ioc, ssl::context& ctx, tcp::endpoint endpoint)
        : ioc_(ioc)
        , ctx_(ctx)
        , acceptor_(net::make_strand(ioc))
    {
        beast::error_code ec;

        // 打开acceptor
        acceptor_.open(endpoint.protocol(), ec);
        if(ec)
        {
            fail(ec, "open");
            return;
        }

        // 设置地址重用
        acceptor_.set_option(net::socket_base::reuse_address(true), ec);
        if(ec)
        {
            fail(ec, "set_option");
            return;
        }

        // 绑定地址
        acceptor_.bind(endpoint, ec);
        if(ec)
        {
            fail(ec, "bind");
            return;
        }

        // 开始监听
        acceptor_.listen(
            net::socket_base::max_listen_connections, ec);
        if(ec)
        {
            fail(ec, "listen");
            return;
        }
    }

    // 开始接受连接
    void run()
    {
        do_accept();
    }

    void do_accept()
    {
        acceptor_.async_accept(
            net::make_strand(ioc_),
            beast::bind_front_handler(
                &listener::on_accept,
                shared_from_this()));
    }

    void on_accept(beast::error_code ec, tcp::socket socket)
    {
        if(ec)
        {
            fail(ec, "accept");
        }
        else
        {
            // 创建新会话并启动
            std::make_shared<session>(std::move(socket), ctx_)->run();
        }

        // 继续接受下一个连接
        do_accept();
    }

    static void fail(beast::error_code ec, char const* what)
    {
        std::cerr << "错误: " << what << ": " << ec.message() << std::endl;
    }
};

int main(int argc, char* argv[])
{
    try
    {
        if(argc != 4)
        {
            std::cerr << "用法: " << argv[0] << " <地址> <端口> <证书文件>" << std::endl;
            return 1;
        }
        auto const address = net::ip::make_address(argv[1]);
        auto const port = static_cast<unsigned short>(std::atoi(argv[2]));
        auto const cert_file = argv[3];

        // 创建IO上下文
        net::io_context ioc{1};

        // 创建SSL上下文
        ssl::context ctx{ssl::context::tlsv12_server};

        // 加载证书文件
        ctx.use_certificate_chain_file(cert_file);
        ctx.use_private_key_file(cert_file, ssl::context::pem);

        // 创建监听器并启动
        std::make_shared<listener>(ioc, ctx, tcp::endpoint{address, port})->run();

        // 运行IO上下文
        ioc.run();
    }
    catch(std::exception const& e)
    {
        std::cerr << "异常: " << e.what() << std::endl;
        return 1;
    }
    return 0;
}

关键修改说明

  1. 异步链式发送:将同步write循环改为do_write异步函数,每次发送完成后通过异步定时器延迟2秒,再发起下一次发送,避免阻塞IO线程。
  2. 恢复读取逻辑:发送完成后重新调用do_read,确保能继续接收客户端消息。
  3. 异常捕获与处理:在fail函数中添加对boost::wrapexcept<boost::system::system_error>的捕获,同时处理broken pipe、stream truncated等特定错误码,避免程序终止。
  4. 缓冲区管理:读取完成后清空缓冲区,避免残留数据影响后续读取。

内容的提问来源于stack exchange,提问作者Rasoul.A

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 06:13:13