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已完全关闭)。
求助问题
- 代码分析中遗漏了什么?
- 是否对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
相关产品推荐
相关产品推荐

