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

如何在Boost::beast同步read()执行时关闭WebSocket连接?

问题:同步read()阻塞时关闭Boost Beast WebSocket连接导致断言错误

简言之:若服务器长时间未发送消息,能否关闭正在执行同步read()操作的WebSocket?

我想用Boost::beast实现一个简单的WebSocket客户端,发现read()是阻塞操作且无法判断是否有消息到来后,创建了一个专门执行read()的线程,即便无数据时阻塞也可接受。

我希望能从非阻塞线程关闭连接,于是调用websocket::close(),但这导致read()抛出BOOST_ASSERT断言错误:

Assertion failed: ! impl.wr_close

请问如何在同步read()执行过程中关闭连接?

复现代码如下:

#include <string>
#include <thread>
#include <chrono>

#include <boost/beast/core.hpp>
#include <boost/beast/websocket.hpp>
#include <boost/asio/connect.hpp>
#include <boost/asio/ip/tcp.hpp>


using namespace std::chrono_literals;

class HandlerThread {

    enum class Status {
        UNINITIALIZED,
        DISCONNECTED,
        CONNECTED,
        READING,
    };

    const std::string _host;
    const std::string _port;
    std::string _resolvedAddress;

    boost::asio::io_context        _ioc;
    boost::asio::ip::tcp::resolver _resolver;
    boost::beast::websocket::stream<boost::asio::ip::tcp::socket> _websocket;

    boost::beast::flat_buffer _buffer;

    bool isRunning = true;
    Status _connectionStatus = Status::UNINITIALIZED;

public:
    HandlerThread(const std::string& host, const uint16_t port)
    : _host(std::move(host))
    , _port(std::to_string(port))
    , _ioc()
    , _resolver(_ioc)
    , _websocket(_ioc) {}

    void Run() {
        // isRunning is also useless, due to blocking boost::beast operations.
        while(isRunning) {
            switch (_connectionStatus) {
                case Status::UNINITIALIZED:
                case Status::DISCONNECTED:
                    if (!connect()) {
                        _connectionStatus = Status::DISCONNECTED;
                        break;
                    }
                case Status::CONNECTED:
                case Status::READING:
                    if (!read()) {
                        _connectionStatus = Status::DISCONNECTED;
                        break;
                    }
            }
        }
    }

    void Close()
    {
         isRunning = false;
        _websocket.close(boost::beast::websocket::close_code::normal);
    }

private:
    bool connect()
    {
        // All here is copy-paste from the examples.
        boost::system::error_code errorCode;
        // Look up the domain name  
        auto const results = _resolver.resolve(_host, _port, errorCode);
        if (errorCode) return false;

        // Make the connection on the IP address we get from a lookup
        auto ep = boost::asio::connect(_websocket.next_layer(), results, errorCode);
        if (errorCode) return false;

        _resolvedAddress = _host + ':' + std::to_string(ep.port());

        _websocket.set_option(boost::beast::websocket::stream_base::decorator(
            [](boost::beast::websocket::request_type& req)
            {
                req.set(boost::beast::http::field::user_agent,
                    std::string(BOOST_BEAST_VERSION_STRING) +
                        " websocket-client-coro");
            }));

        boost::beast::websocket::response_type res;
        _websocket.handshake(res, _resolvedAddress, "/", errorCode);

        if (errorCode) return false;

        _connectionStatus = Status::CONNECTED;
        return true;
    }

    bool read()
    {
        boost::system::error_code errorCode;
        _websocket.read(_buffer, errorCode);

        if (errorCode) return false;

        if (_websocket.is_message_done()) {
            _connectionStatus = Status::CONNECTED;
            // notifyRead(_buffer);
            _buffer.clear();    
        } else {
            _connectionStatus = Status::READING;
        }

        return true;
    }
};

int main() {
    HandlerThread handler("localhost", 8080);
    std::thread([&]{
        handler.Run();
    }).detach(); // bye!

    std::this_thread::sleep_for(3s);
    handler.Close(); // Bad idea...

    return 0;
}

解决方案

核心原因

Boost Beast的WebSocket对象不支持线程安全操作,跨线程同时调用read()和close()会触发内部状态竞争,直接导致断言失败。所有WebSocket操作必须在其所属io_context的运行线程中执行。

修正步骤

  1. 通过io_context投递关闭操作
    不在外部线程直接调用websocket::close(),而是用boost::asio::post将关闭任务投递到WebSocket所属的io_context中,让任务在Run()函数所在线程执行,避免线程竞争。

  2. 修复线程安全的状态变量
    将isRunning改为std::atomic<bool>,确保多线程访问时的内存可见性,避免未定义行为。

  3. 处理read操作的错误
    当read()返回错误时,判断是否为关闭相关错误(如operation_aborted或closed),并设置isRunning为false,让循环正常退出。

修改后的完整代码

#include <string>
#include <thread>
#include <chrono>
#include <atomic>

#include <boost/beast/core.hpp>
#include <boost/beast/websocket.hpp>
#include <boost/asio/connect.hpp>
#include <boost/asio/ip/tcp.hpp>


using namespace std::chrono_literals;

class HandlerThread {

    enum class Status {
        UNINITIALIZED,
        DISCONNECTED,
        CONNECTED,
        READING,
    };

    const std::string _host;
    const std::string _port;
    std::string _resolvedAddress;

    boost::asio::io_context        _ioc;
    boost::asio::ip::tcp::resolver _resolver;
    boost::beast::websocket::stream<boost::asio::ip::tcp::socket> _websocket;

    boost::beast::flat_buffer _buffer;

    std::atomic<bool> isRunning = true;
    Status _connectionStatus = Status::UNINITIALIZED;

public:
    HandlerThread(const std::string& host, const uint16_t port)
    : _host(std::move(host))
    , _port(std::to_string(port))
    , _ioc()
    , _resolver(_ioc)
    , _websocket(_ioc) {}

    void Run() {
        while(isRunning) {
            switch (_connectionStatus) {
                case Status::UNINITIALIZED:
                case Status::DISCONNECTED:
                    if (!connect()) {
                        _connectionStatus = Status::DISCONNECTED;
                        break;
                    }
                case Status::CONNECTED:
                case Status::READING:
                    if (!read()) {
                        _connectionStatus = Status::DISCONNECTED;
                        break;
                    }
            }
        }
    }

    void Close()
    {
         isRunning = false;
         // 通过io_context投递关闭操作到Run线程执行
         boost::asio::post(_ioc, [this]() {
             boost::system::error_code ec;
             _websocket.close(boost::beast::websocket::close_code::normal, ec);
             // 忽略关闭时的非致命错误,比如连接已断开
         });
    }

private:
    bool connect()
    {
        boost::system::error_code errorCode;
        auto const results = _resolver.resolve(_host, _port, errorCode);
        if (errorCode) return false;

        auto ep = boost::asio::connect(_websocket.next_layer(), results, errorCode);
        if (errorCode) return false;

        _resolvedAddress = _host + ':' + std::to_string(ep.port());

        _websocket.set_option(boost::beast::websocket::stream_base::decorator(
            [](boost::beast::websocket::request_type& req)
            {
                req.set(boost::beast::http::field::user_agent,
                    std::string(BOOST_BEAST_VERSION_STRING) +
                        " websocket-client-coro");
            }));

        boost::beast::websocket::response_type res;
        _websocket.handshake(res, _resolvedAddress, "/", errorCode);

        if (errorCode) return false;

        _connectionStatus = Status::CONNECTED;
        return true;
    }

    bool read()
    {
        boost::system::error_code errorCode;
        _websocket.read(_buffer, errorCode);

        if (errorCode) {
            // 处理关闭相关错误,终止循环
            if (errorCode == boost::asio::error::operation_aborted || 
                errorCode == boost::beast::websocket::error::closed) {
                isRunning = false;
            }
            return false;
        }

        if (_websocket.is_message_done()) {
            _connectionStatus = Status::CONNECTED;
            // notifyRead(_buffer);
            _buffer.clear();    
        } else {
            _connectionStatus = Status::READING;
        }

        return true;
    }
};

int main() {
    HandlerThread handler("localhost", 8080);
    std::thread worker([&]{
        handler.Run();
    });

    std::this_thread::sleep_for(3s);
    handler.Close();

    worker.join(); // 等待线程退出,避免资源泄漏
    return 0;
}

额外优化建议

  • 避免使用detach()线程,改用join()等待线程正常退出,防止程序提前结束导致资源泄漏。
  • 可以给read操作设置超时,通过websocket::set_option配置stream_base::timeout,避免无限期阻塞。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 11:34:57