如何在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的运行线程中执行。
修正步骤
通过io_context投递关闭操作
不在外部线程直接调用websocket::close(),而是用boost::asio::post将关闭任务投递到WebSocket所属的io_context中,让任务在Run()函数所在线程执行,避免线程竞争。修复线程安全的状态变量
将isRunning改为std::atomic<bool>,确保多线程访问时的内存可见性,避免未定义行为。处理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
相关产品推荐
相关产品推荐

