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

Boost.Beast异步WebSocket客户端同实例重连失败问题求助

Boost.Beast WebSocket 同实例重连失败问题解决

问题根源

复用同一个 websocket::stream 实例时,若未彻底清理前一次连接的状态,会导致底层 SSL 流残留旧连接的上下文(如未清空的缓冲区、未重置的 SSL 会话状态)。即使 TCP 重连成功,SSL 握手或 WebSocket 握手阶段会因收到不符合预期的消息触发 ssl3_read_bytes 错误。

可行解决方案:同实例重连的正确步骤

要复用 session 实例,需按以下流程彻底清理并重置连接状态:

1. 优雅清理现有连接

在重连前,必须终止所有异步操作并逐层关闭连接:

boost::system::error_code ec;

// 取消所有未完成的异步操作
ws_->cancel(ec);

// 尝试正常关闭 WebSocket 连接(忽略已关闭或被中止的错误)
ws_->close(websocket::close_code::normal, ec);
if (ec && ec != beast::error::closed && ec != boost::asio::error::operation_aborted) {
    std::cerr << "WebSocket 关闭错误: " << ec.message() << std::endl;
}

// 关闭 SSL 流(忽略 EOF 或中止错误)
ws_->next_layer().shutdown(ssl::stream_base::shutdown_both, ec);
if (ec && ec != boost::asio::error::eof && ec != boost::asio::error::operation_aborted) {
    std::cerr << "SSL 关闭错误: " << ec.message() << std::endl;
}

// 关闭底层 TCP 套接字
ws_->next_layer().next_layer().close(ec);
if (ec && ec != boost::asio::error::not_open) {
    std::cerr << "TCP 关闭错误: " << ec.message() << std::endl;
}

2. 重置 WebSocket 流

使用智能指针(如 std::unique_ptr)封装 WebSocket 流,重连时直接创建新实例以彻底清除旧状态:

// 假设 session 类持有 io_context 和 ssl::context 的引用
ws_ = std::make_unique<websocket::stream<ssl::stream<tcp::socket>>>(ioc_, ctx_);

3. 重新执行连接流程

重置完成后,重新执行 DNS 解析、TCP 连接、SSL 握手、WebSocket 握手的完整流程,与首次连接逻辑一致。

完整修改后的 Session 示例代码

以下是集成重连逻辑的 session 类实现,包含消息队列发送功能:

#include <boost/beast.hpp>
#include <boost/asio.hpp>
#include <boost/asio/ssl.hpp>
#include <memory>
#include <mutex>
#include <deque>
#include <iostream>

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

class session : public std::enable_shared_from_this<session> {
public:
    session(net::io_context& ioc, ssl::context& ctx, std::string host, std::string port)
        : ioc_(ioc), ctx_(ctx), host_(std::move(host)), port_(std::move(port)),
          ws_(std::make_unique<websocket::stream<ssl::stream<tcp::socket>>>(ioc, ctx)) {}

    void start() {
        do_resolve();
    }

    // 触发重连
    void reconnect() {
        boost::system::error_code ec;

        // 清理现有连接
        ws_->cancel(ec);
        ws_->close(websocket::close_code::normal, ec);
        ws_->next_layer().shutdown(ssl::stream_base::shutdown_both, ec);
        ws_->next_layer().next_layer().close(ec);

        // 重置 WebSocket 流
        ws_ = std::make_unique<websocket::stream<ssl::stream<tcp::socket>>>(ioc_, ctx_);

        // 重新启动连接流程
        do_resolve();
    }

    // 向队列添加待发送消息
    void send(std::string msg) {
        std::lock_guard<std::mutex> lock(mutex_);
        queue_.push_back(std::move(msg));
        if (!writing_) {
            do_write();
        }
    }

private:
    void do_resolve() {
        auto self(shared_from_this());
        tcp::resolver resolver(ioc_);
        resolver.async_resolve(host_, port_,
            [self](boost::system::error_code ec, tcp::resolver::results_type results) {
                if (!ec) {
                    self->do_connect(results);
                } else {
                    std::cerr << "DNS 解析错误: " << ec.message() << std::endl;
                    net::post(self->ioc_, [self]() { self->reconnect(); });
                }
            });
    }

    void do_connect(tcp::resolver::results_type const& results) {
        auto self(shared_from_this());
        net::async_connect(ws_->next_layer().next_layer(), results,
            [self](boost::system::error_code ec, tcp::endpoint const&) {
                if (!ec) {
                    self->do_ssl_handshake();
                } else {
                    std::cerr << "TCP 连接错误: " << ec.message() << std::endl;
                    net::post(self->ioc_, [self]() { self->reconnect(); });
                }
            });
    }

    void do_ssl_handshake() {
        auto self(shared_from_this());
        ws_->next_layer().async_handshake(ssl::stream_base::client,
            [self](boost::system::error_code ec) {
                if (!ec) {
                    self->do_ws_handshake();
                } else {
                    std::cerr << "SSL 握手错误: " << ec.message() << std::endl;
                    net::post(self->ioc_, [self]() { self->reconnect(); });
                }
            });
    }

    void do_ws_handshake() {
        auto self(shared_from_this());
        ws_->async_handshake(host_, "/",
            [self](boost::system::error_code ec) {
                if (!ec) {
                    std::cout << "连接成功" << std::endl;
                    self->do_read();
                    // 发送队列中积压的消息
                    std::lock_guard<std::mutex> lock(self->mutex_);
                    if (!self->queue_.empty()) {
                        self->do_write();
                    }
                } else {
                    std::cerr << "WebSocket 握手错误: " << ec.message() << std::endl;
                    net::post(self->ioc_, [self]() { self->reconnect(); });
                }
            });
    }

    void do_read() {
        auto self(shared_from_this());
        ws_->async_read(buffer_,
            [self](boost::system::error_code ec, std::size_t bytes_transferred) {
                boost::ignore_unused(bytes_transferred);
                if (!ec) {
                    std::cout << "收到消息: " << beast::buffers_to_string(self->buffer_.data()) << std::endl;
                    self->buffer_.consume(bytes_transferred);
                    self->do_read();
                } else {
                    std::cerr << "读取错误: " << ec.message() << std::endl;
                    self->reconnect();
                }
            });
    }

    void do_write() {
        auto self(shared_from_this());
        std::lock_guard<std::mutex> lock(mutex_);
        if (queue_.empty()) {
            writing_ = false;
            return;
        }
        writing_ = true;
        auto msg = std::move(queue_.front());
        queue_.pop_front();
        ws_->async_write(net::buffer(msg),
            [self, msg = std::move(msg)](boost::system::error_code ec, std::size_t bytes_transferred) {
                boost::ignore_unused(bytes_transferred);
                if (!ec) {
                    std::lock_guard<std::mutex> lock(self->mutex_);
                    self->do_write();
                } else {
                    std::cerr << "发送错误: " << ec.message() << std::endl;
                    // 将未发送的消息重新放回队列
                    std::lock_guard<std::mutex> lock(self->mutex_);
                    self->queue_.push_front(std::move(msg));
                    self->writing_ = false;
                    self->reconnect();
                }
            });
    }

    net::io_context& ioc_;
    ssl::context& ctx_;
    std::string host_;
    std::string port_;
    std::unique_ptr<websocket::stream<ssl::stream<tcp::socket>>> ws_;
    beast::flat_buffer buffer_;
    std::mutex mutex_;
    std::deque<std::string> queue_;
    bool writing_ = false;
};

关键说明

  • 使用 std::unique_ptr 封装 WebSocket 流,确保重连时彻底销毁旧实例并创建全新流,避免状态残留。
  • 重连时保留消息队列,连接恢复后自动发送积压消息。
  • 所有错误分支均触发重连逻辑,保证客户端的连接稳定性。

内容的提问来源于stack exchange,提问作者darnell_a

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 22:52:34