如何将TCP Socket接收的数据转发至WebSocket客户端?程序崩溃排查
问题分析与修复
崩溃核心原因
- 全局变量竞态:
is_websocket_initialized全局变量未加锁,多线程环境下读写会导致session指针未初始化就被tcp_session引用,或状态判断逻辑混乱。 - 缓冲区生命周期失效:tcp_session的buffer在发起
async_write后立即被consume,但异步写操作尚未完成,后续会访问已失效的缓冲区内存。 - Strand线程安全违规:tcp_session在自身strand线程中直接调用session的
async_write,违反Boost Asio规则——WebSocket流的异步操作必须在其绑定的strand上发起,跨strand调用会引发线程安全问题。 - 缓冲区重复提交:tcp_session的
on_read中两次调用buffer_.commit(bytes_transferred),导致缓冲区数据状态异常。
修复步骤
1. 修复全局变量线程安全
新增互斥锁保护全局session指针和初始化状态:
std::mutex g_session_mutex; std::shared_ptr<session> g_global_session; bool is_websocket_initialized = false;
读写这些变量时加锁,例如listener创建session时:
std::lock_guard<std::mutex> lock(g_session_mutex); if (!is_websocket_initialized) { is_websocket_initialized = true; g_global_session = std::make_shared<session>(std::move(socket)); g_global_session->run(); }
2. 确保缓冲区生命周期覆盖异步操作
将tcp_session的缓冲区数据转为字符串后传入session,避免异步操作访问已被清理的缓冲区:
// tcp_session的on_read函数中 buffer_.commit(bytes_transferred); std::string data = beast::buffers_to_string(buffer_.data()); buffer_.consume(buffer_.size()); std::lock_guard<std::mutex> lock(g_session_mutex); if (g_global_session) { g_global_session->write_data_to_websocket(data); }
3. 强制在WebSocket的strand上发起异步操作
使用net::dispatch将写操作切换到WebSocket流的executor(即strand)上,保证线程安全:
// session的write_data_to_websocket函数中 net::dispatch(ws_.get_executor(), [self = shared_from_this()]() { self->ws_.text(true); self->ws_.async_write( self->buffer_.data(), beast::bind_front_handler( &session::on_write, self)); });
4. 移除重复的缓冲区commit
删除tcp_session::on_read中重复的buffer_.commit(bytes_transferred)调用。
完整修复代码
#include <boost/beast/core.hpp> #include <boost/beast/websocket.hpp> #include <boost/asio/dispatch.hpp> #include <boost/asio/strand.hpp> #include <algorithm> #include <cstdlib> #include <functional> #include <iostream> #include <memory> #include <string> #include <thread> #include <vector> #include <mutex> namespace beast = boost::beast; // from <boost/beast.hpp> namespace http = beast::http; // from <boost/beast/http.hpp> namespace websocket = beast::websocket; // from <boost/beast/websocket.hpp> namespace net = boost::asio; // from <boost/asio.hpp> using tcp = boost::asio::ip::tcp; // from <boost/asio/ip/tcp.hpp> //------------------------------------------------------------------------------ std::mutex g_session_mutex; std::shared_ptr<session> g_global_session; bool is_websocket_initialized = false; void fail(beast::error_code ec, char const* what) { std::cerr << what << ": " << ec.message() << "\n"; } class session : public std::enable_shared_from_this<session> { websocket::stream<beast::tcp_stream> ws_; beast::flat_buffer buffer_; public: explicit session(tcp::socket&& socket) : ws_(std::move(socket)) { } void run() { net::dispatch(ws_.get_executor(), beast::bind_front_handler( &session::on_run, shared_from_this())); } void on_run() { ws_.set_option( websocket::stream_base::timeout::suggested( beast::role_type::server)); ws_.set_option(websocket::stream_base::decorator( [](websocket::response_type& res) { res.set(http::field::server, std::string(BOOST_BEAST_VERSION_STRING) + " websocket-server-async"); })); 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"); std::cout << "WebSocket connection initialized" << std::endl; std::lock_guard<std::mutex> lock(g_session_mutex); is_websocket_initialized = true; 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) { std::lock_guard<std::mutex> lock(g_session_mutex); is_websocket_initialized = false; g_global_session.reset(); return; } if (ec) return fail(ec, "read"); ws_.text(ws_.got_text()); ws_.async_write( buffer_.data(), beast::bind_front_handler( &session::on_write, shared_from_this())); } void write_data_to_websocket(const std::string& data) { std::cout << "received : " << data << "\n"; beast::ostream(buffer_) << data; net::dispatch(ws_.get_executor(), [self = shared_from_this()]() { self->ws_.text(true); self->ws_.async_write( self->buffer_.data(), beast::bind_front_handler( &session::on_write, self)); }); } void on_write( beast::error_code ec, std::size_t bytes_transferred) { boost::ignore_unused(bytes_transferred); if (ec) return fail(ec, "write"); buffer_.consume(buffer_.size()); do_read(); } }; class tcp_session : public std::enable_shared_from_this<tcp_session> { tcp::socket socket_; beast::flat_buffer buffer_; public: explicit tcp_session(tcp::socket&& socket) : socket_(std::move(socket)) {} void run() { do_read(); } void do_read() { auto self = shared_from_this(); socket_.async_read_some(buffer_.prepare(1024), [self](beast::error_code ec, std::size_t bytes_transferred) { self->on_read(ec, bytes_transferred); }); } void on_read(beast::error_code ec, std::size_t bytes_transferred) { if (ec) { fail(ec, "tcp read"); return; } buffer_.commit(bytes_transferred); std::string data = beast::buffers_to_string(buffer_.data()); buffer_.consume(buffer_.size()); std::lock_guard<std::mutex> lock(g_session_mutex); if (g_global_session) { g_global_session->write_data_to_websocket(data); } do_read(); } }; class listener : public std::enable_shared_from_this<listener> { net::io_context& ioc_; tcp::acceptor acceptor_; bool is_websocket; public: listener( net::io_context& ioc, tcp::endpoint endpoint, bool ws_flag) : ioc_(ioc) , acceptor_(ioc) , is_websocket(ws_flag) { beast::error_code ec; 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(); } private: 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 if (is_websocket) { std::lock_guard<std::mutex> lock(g_session_mutex); if (!is_websocket_initialized) { std::cout << "creating websocket" << "\n"; g_global_session = std::make_shared<session>(std::move(socket)); g_global_session->run(); } else { std::cout << "WebSocket already initialized, rejecting new connection" << "\n"; } } else { std::cout << "receiving as tcp_socket" << "\n"; auto tcpSess = std::make_shared<tcp_session>(std::move(socket)); tcpSess->run(); } do_accept(); } }; //------------------------------------------------------------------------------ int main(int argc, char* argv[]) { if (argc != 4) { std::cerr << "Usage: websocket-server-async <address> <port> <threads>\n" << "Example:\n" << " websocket-server-async 0.0.0.0 8080 1\n"; return EXIT_FAILURE; } auto const address = net::ip::make_address(argv[1]); auto const port = static_cast<unsigned short>(std::atoi(argv[2])); auto const threads = std::max<int>(1, std::atoi(argv[3])); net::io_context ioc{ threads }; auto const test_port = 7777; std::make_shared<listener>(ioc, tcp::endpoint{ address, port }, true)->run(); std::make_shared<listener>(ioc, tcp::endpoint{ address, test_port }, false)->run(); std::vector<std::thread> v; v.reserve(threads - 1); for (auto i = threads - 1; i > 0; --i) v.emplace_back( [&ioc] { ioc.run(); }); ioc.run(); return EXIT_SUCCESS; }
内容的提问来源于stack exchange,提问作者UVision
相关产品推荐
相关产品推荐

