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

Boost Beast WebSocket服务器连接关闭后堆内存未完全释放问题排查

Boost Beast WebSocket 堆内存持续增长问题排查求助

基于Boost官方异步服务器示例修改的WebSocket服务器,出现无法解释的堆内存分配问题。使用_CrtMemState和_CrtMemCheckpoint监控堆内存,测试场景为客户端主动发起连接并关闭,每次仅处理单个连接,通过Python客户端多次连接断开(每次保持数秒,断开后等待至少2秒重复),每500ms记录内存状态,发现会话销毁后内存未回退:

  • 初始启动:25009(就绪状态)
  • 首次断开后:29752
  • 第二次断开后:30688
  • 第三次断开后:30873

首次断开后的内存增长推测来自静态对象构造,但后续每次断开后内存应回落至首次断开后的数值,实际却持续上升。

服务器架构:WebSocketManager实例化Listener并运行于io_context,Session对象注册到管理器,由独立线程检测状态并销毁已关闭会话。已尝试在Session析构中关闭Socket、清理缓冲区,但仍存在内存残留,调用drain时出现10009错误(Socket已完全关闭)。

求助问题

  1. 代码分析中遗漏了什么?
  2. 是否对WebSocket或缓冲区的关闭/清理逻辑存在误解?

完整可复现代码

//
// Copyright (c) 2017 Vinnie Falco (vinnie dot falco at gmail dot com)
//
// Distributed under the Boost Software License, Version 1.0. (See accompanying
// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
//
// Official repository: https://github.com/boostorg/beast
//

//#include "../common/helpers.hpp"
#define _CRTDBG_MAP_ALLOC
#include <boost/beast/core.hpp>
#include <boost/beast/websocket.hpp>
#include <boost/asio.hpp>
#include <boost/optional.hpp>
#include <chrono>
#include <cstdlib>
#include <ctime>
#include <iostream>
#include <memory>
#include <string>
#include <thread>

#include <crtdbg.h>

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>

static void socket_fail(beast::error_code ec, char const* what)
{
    std::cerr << what << ": " << ec.message() << "\n";
}

class Session
    : public std::enable_shared_from_this<Session>
{
private:
    websocket::stream<beast::tcp_stream> ws_;
    //beast::flat_buffer buffer_out_;
    beast::flat_buffer buffer_in_;

public:
    Session(tcp::socket&& socket)
        : ws_(std::move(socket))
    {

    }

    virtual ~Session()
    {
        std::cout << "Session Destroyed" << std::endl;
        boost::beast::error_code ec;
        buffer_in_.consume(buffer_in_.size() + 1);
        ws_.next_layer().cancel();
        ws_.next_layer().socket().lowest_layer().close();
        ws_.close(boost::beast::websocket::close_code::normal, ec);
    }

    void run()
    {
        // We need to be executing within a strand to perform async operations
        // on the I/O objects in this session. Although not strictly necessary
        // for single-threaded contexts, this example code is written to be
        // thread-safe by default.

        net::dispatch(ws_.get_executor(),
            beast::bind_front_handler(
                &Session::on_run,
                shared_from_this()));
    }

    void on_run()
    {
        // Set suggested timeout settings for the websocket
        ws_.set_option(
            websocket::stream_base::timeout::suggested(
                beast::role_type::server));

        // Set a decorator to change the Server of the handshake
        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_.read_message_max(64 * 1024 * 1024);
        //ws_.write_buffer_bytes(16 * 1024 * 1024);

        // Accept the websocket handshake
        ws_.async_accept(
            beast::bind_front_handler(
                &Session::on_accept,
                shared_from_this()));
    }

    void on_accept(beast::error_code ec)
    {
        if (ec)
            return socket_fail(ec, "accept");

        // Read a message
        ws_.async_read(
            buffer_in_,
            beast::bind_front_handler(
                &Session::on_read,
                shared_from_this()));
    }

    void on_read(beast::error_code ec, std::size_t bytes_transferred)
    {
        // This indicates that the session was closed
        if (ec == websocket::error::closed) {
            return;
        }
        else if (ec == beast::error::timeout)
        {
            std::cout << "Read Request Failed" << std::endl;
        }
        else if (ec) {
            socket_fail(ec, "read");
        }

        buffer_in_.consume(buffer_in_.size());
        buffer_in_.clear();
        if (ws_.is_open())
        {
            ws_.async_read(
                buffer_in_,
                beast::bind_front_handler(
                    &Session::on_read,
                    shared_from_this()));
        }
    }

    bool isRunning()
    {
        return ws_.is_open();
    }
};



class WebSocketManager
{

private:
    WebSocketManager()
    {
        __session_update_timer = std::thread(&WebSocketManager::updateSessions, this);
    }

    static WebSocketManager* __instance;

    ~WebSocketManager()
    {

    }

    std::thread __socket_thread;

    void run();

    bool __is_running = false;

    std::vector<std::shared_ptr<Session>> __sessions;
    std::thread __session_update_timer;

    void updateSessions()
    {
        while (true)
        {
            using namespace std::chrono_literals;
            std::weak_ptr<Session> weak_session;
            for (std::vector<std::shared_ptr<Session>>::iterator itr = __sessions.begin(); itr != __sessions.end();)
            {
                if (itr->get()->isRunning())
                {
                    ++itr;
                }
                else
                {
                    weak_session = *itr;
                    //(*itr).reset();
                    itr = __sessions.erase(itr);
                }
            }
            std::this_thread::sleep_for(500ms);
        }
    }

protected:
    void registerSession(std::shared_ptr<Session>& pSession)
    {
        __sessions.push_back(pSession);
    }

public:
    WebSocketManager(WebSocketManager const&) = delete;
    void operator=(WebSocketManager const&) = delete;

    static WebSocketManager* Instance()
    {
        if (__instance == nullptr)
        {
            __instance = new WebSocketManager();
        }
        return __instance;
    }

    void startBehavior()
    {
        if (!__is_running)
        {
            __is_running = true;

            __socket_thread = std::thread(std::bind(&WebSocketManager::run, this));
        }
    }

    size_t getActiveSessions() const
    {
        return __sessions.size();
    }

    const std::string SRV_ADDRESS = "127.0.0.1";
    const int PORT = 617;
    const int MAX_THREAD_COUNT = 5;

    friend class Listener;
};

class Listener
    : public std::enable_shared_from_this<Listener>
{
private:
    net::io_context& ioc_;
    tcp::acceptor acceptor_;

    Listener();

public:
    Listener(net::io_context& ioc, tcp::endpoint endpoint)
        :
        ioc_(ioc)
        , acceptor_(ioc)
    {
        beast::error_code ec;

        // Open the acceptor
        acceptor_.open(endpoint.protocol(), ec);
        if (ec)
        {
            socket_fail(ec, "open");
            return;
        }

        // Allow address reuse
        acceptor_.set_option(net::socket_base::reuse_address(true), ec);
        if (ec)
        {
            socket_fail(ec, "set_option");
            return;
        }

        // Bind to the server address
        acceptor_.bind(endpoint, ec);
        if (ec)
        {
            socket_fail(ec, "bind");
            return;
        }

        // Start listening for connections
        acceptor_.listen(
            net::socket_base::max_listen_connections, ec);
        if (ec)
        {
            socket_fail(ec, "listen");
            return;
        }
    }

    void run()
    {
        do_accept();
    }

    void do_accept()
    {
        // The new connection gets its own strand
        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)
        {
            socket_fail(ec, "accept");
        }
        else
        {
            std::cout << "Session create" << std::endl;
            // Create the session and run it
            std::shared_ptr<Session> session = std::make_shared<Session>(std::move(socket));
            WebSocketManager::Instance()->registerSession(session);
            session->run();
        }

        // Accept another connection
        do_accept();
    }
};

WebSocketManager* WebSocketManager::__instance(nullptr);

void WebSocketManager::run()
{
    auto const address = net::ip::make_address(SRV_ADDRESS.c_str());
    auto const threads = std::max(1, 1);

    net::io_context ioc{ threads };

    std::make_shared<Listener>(ioc, net::ip::tcp::endpoint{ address, static_cast<unsigned short>(PORT) })->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();
}



int main(int argc, char* argv[])
{
    using namespace std::chrono_literals;
    WebSocketManager::Instance()->startBehavior();
    _CrtMemState S1;
    uint16_t counter = 600;
    bool session_detected = false;
    int skip = true;

    _CrtMemCheckpoint(&S1);

    while (counter > 0)
    {
        _CrtMemCheckpoint(&S1);
        std::cout << S1.lSizes[1] << "," << S1.lSizes[2] << std::endl;
        if (WebSocketManager::Instance()->getActiveSessions() > 0)
        {
            session_detected = true;
        }
        else if (WebSocketManager::Instance()->getActiveSessions() == 0 && session_detected)
        {
            if (skip) skip = false;
            else
            {
                session_detected = false;
                _CrtDumpMemoryLeaks();
                skip = true;
            }
        }
        std::this_thread::sleep_for(500ms);
    }

    return EXIT_SUCCESS;
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 04:07:33