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

如何将TCP Socket接收的数据转发至WebSocket客户端?程序崩溃排查

问题分析与修复

崩溃核心原因

  1. 全局变量竞态:is_websocket_initialized全局变量未加锁,多线程环境下读写会导致session指针未初始化就被tcp_session引用,或状态判断逻辑混乱。
  2. 缓冲区生命周期失效:tcp_session的buffer在发起async_write后立即被consume,但异步写操作尚未完成,后续会访问已失效的缓冲区内存。
  3. Strand线程安全违规:tcp_session在自身strand线程中直接调用session的async_write,违反Boost Asio规则——WebSocket流的异步操作必须在其绑定的strand上发起,跨strand调用会引发线程安全问题。
  4. 缓冲区重复提交: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 17:40:53