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

调用boost::concurrent_flat_map的erase方法导致程序崩溃问题排查

修复Boost UDP服务器中concurrent_flat_map erase导致的崩溃问题

问题诊断

崩溃并非直接由boost::concurrent_flat_map::erase导致,而是多线程环境下的内存管理、线程生命周期及回调捕获错误引发的未定义行为,具体问题包括:

  • 栈内存悬空:在UDP接收回调中使用栈数组buffer存储消息,将指针传入UdpClient的队列后,栈帧销毁导致指针变为野指针,后续队列处理或销毁时delete该指针会触发内存错误。
  • 线程未正确回收:UdpClient的处理线程m_proc_thread未在析构时执行join,当shared_ptr释放UdpClient时,线程可能仍在运行,访问已释放的成员变量(如m_on_disconnected)。
  • 回调捕获风险:创建UdpClient时的断开回调使用[&]捕获,若服务器实例提前销毁,会导致回调访问无效内存;同时回调执行时可能与handle_receive中的visit操作并发,虽concurrent_flat_map线程安全,但野指针问题已埋下崩溃隐患。

修复方案

1. 替换栈内存为堆内存

接收消息时动态分配内存,确保指针在队列处理周期内有效:

// 原代码
unsigned char buffer[num_bytes];
std::copy_n(std::begin(this->m_buffer), num_bytes, buffer);

// 修改为
unsigned char* buffer = new unsigned char[num_bytes];
std::copy_n(std::begin(this->m_buffer), num_bytes, buffer);

2. 为UdpClient添加析构函数,回收线程

确保线程完成所有操作后再销毁对象:

~UdpClient() {
    m_running = false;
    m_proc_cond.notify_one();
    if (m_proc_thread.joinable()) {
        m_proc_thread.join();
    }
}

3. 安全捕获服务器实例

让UdpServer继承std::enable_shared_from_this,使用shared_ptr捕获服务器实例,避免回调访问无效对象:

class UdpServer : public std::enable_shared_from_this<UdpServer> {
    // ... 其他代码
};

// 创建UdpClient时的回调修改为
auto self = shared_from_this();
auto client = std::make_shared<UdpClient>(recv_id,
    [self, id]() {
        std::cout << "received disconnect signal\n";
        if (self->m_clients.contains(id)) {        
            self->m_clients.erase(id);
        }      
    }
);

4. 修正条件变量的等待逻辑

在UdpClient::handle_queue中,添加原子标志控制线程退出,避免线程在对象销毁后继续运行:

// UdpClient添加成员
std::atomic<bool> m_running{true};

// handle_queue修改为
void handle_queue() {
    while (m_running) {       
        while (this->m_queue.read_available() > 0 && m_running)            
        {
            this->m_queue.consume_one(
                [&] (unsigned char* buffer)
                {
                    // Process message here
                    delete[] buffer; // 处理完后释放内存
                }
            );
        } 
      
        std::unique_lock<std::mutex> lk(this->m_proc_mutex);
        auto now     = std::chrono::system_clock::now();
        auto timeout = now + std::chrono::seconds(5);
        if (!this->m_proc_cond.wait_until(lk, timeout, [&] ()
        {
            return !m_running || this->m_queue.read_available() > 0;
        })) 
        {   
            std::cout << "timeout occurs\n";
            break;
        }
    } 
    
    // Cleanup messages queue
    this->m_queue.consume_all([] (unsigned char* buffer) { delete[] buffer; });

    // Notify disconnection only if we were running
    if (m_running) {
        m_on_disconnected(this->m_recv_id);
    }
}

修正后的完整代码

#include <stdio.h>
#include <iostream>
#include <string>
#include <chrono>
#include <condition_variable>
#include <atomic>
#include <unistd.h>  
#include <thread>
#include <functional>
#include <memory>

#include <boost/asio.hpp>
#include <boost/thread.hpp>
#include <boost/unordered/concurrent_flat_map.hpp>
#include <boost/lockfree/spsc_queue.hpp>


class UdpClient
{
    public:
        // Constructor
        UdpClient(int recv_id, std::function<void(int)> cb) :
            m_queue(100),
            m_running(true)
        {
            this->m_recv_id = recv_id;
            this->m_on_disconnected = cb;
            this->m_proc_thread = std::thread(std::bind(&UdpClient::handle_queue, this));
        }

        // Destructor
        ~UdpClient() {
            m_running = false;
            m_proc_cond.notify_one();
            if (m_proc_thread.joinable()) {
                m_proc_thread.join();
            }
        }

        // Enqueue
        void enqueue(unsigned char* buffer) 
        {   
            if (m_running && this->m_queue.write_available() > 0)
            {
                this->m_queue.push(buffer);
                this->m_proc_cond.notify_one();
            } else {
                delete[] buffer; // 若队列不可用,直接释放内存
            }
        }

    private:

        int m_recv_id;    
        std::function<void(int)> m_on_disconnected;
        boost::lockfree::spsc_queue<unsigned char*> m_queue;
        std::condition_variable m_proc_cond;
        std::mutex              m_proc_mutex;
        std::thread             m_proc_thread;
        std::atomic<bool>       m_running;

        // Handle queue
        void handle_queue()
        {
            while (m_running) 
            {       
                while (this->m_queue.read_available() > 0 && m_running)            
                {
                    this->m_queue.consume_one(
                        [&] (unsigned char* buffer)
                        {
                            // Process message here
                            delete[] buffer; // 释放消息内存
                        }
                    );
                } 
              
                std::unique_lock<std::mutex> lk(this->m_proc_mutex);
                auto now     = std::chrono::system_clock::now();
                auto timeout = now + std::chrono::seconds(5);
                if (!this->m_proc_cond.wait_until(lk, timeout, [&] ()
                {
                    return !m_running || this->m_queue.read_available() > 0;
                })) 
                {   
                    std::cout << "timeout occurs\n";
                    break;
                }
            } 
            
            // Cleanup messages queue
            this->m_queue.consume_all([] (unsigned char* buffer) { delete[] buffer; });

            // Notify disconnection only if we were running
            if (m_running) {
                m_on_disconnected(this->m_recv_id);
            }
        }

};

class UdpServer : public std::enable_shared_from_this<UdpServer>
{
    public:
        // Constructor
        UdpServer(int port) : 
            m_ioservice(),    
            m_clients  (5)
        {
            this->m_port = port;
            this->m_socket = std::make_shared<boost::asio::ip::udp::socket>(m_ioservice);
        }

        // Listen
        void listen()
        {
            auto endpoint = boost::asio::ip::udp::endpoint(boost::asio::ip::udp::v4(), this->m_port);
            this->m_socket->open(endpoint.protocol());
            this->m_socket->set_option(boost::asio::ip::udp::socket::reuse_address(true));
            this->m_socket->set_option(boost::asio::detail::socket_option::boolean<SOL_SOCKET, SO_REUSEPORT>(true));
            this->m_socket->set_option(boost::asio::ip::udp::socket::keep_alive(true));
            this->m_socket->bind(endpoint);
            this->handle_receive();
            this->m_ioservice.run();
        }
    
    private:

        int m_port;
        boost::asio::io_service m_ioservice;    
        std::shared_ptr<boost::asio::ip::udp::socket> m_socket;   
        boost::concurrent_flat_map<int32_t, std::shared_ptr<UdpClient>> m_clients; 
        std::array<unsigned char, 1500> m_buffer {};
        boost::asio::ip::udp::endpoint m_remote_endpoint;

        // Handle receive
        void handle_receive() 
        {  
            this->m_socket->async_receive_from(
                boost::asio::buffer(this->m_buffer, 1500), 
                this->m_remote_endpoint,
                [this] (boost::system::error_code const & err, size_t num_bytes)
                {    
                    if (!err && num_bytes > 0) 
                    {                
                        int recv_id = 1;
                        // 使用堆内存存储消息
                        unsigned char* buffer = new unsigned char[num_bytes];
                        std::copy_n(std::begin(this->m_buffer), num_bytes, buffer);
                                                
                        if (this->m_clients.contains(recv_id))
                        {
                            this->m_clients.visit(recv_id,
                                [&] (std::pair<const int, std::shared_ptr<UdpClient>> &x)
                                {
                                    x.second->enqueue(buffer);
                                }
                            );                            
                        }
                        else
                        {
                            auto self = shared_from_this();
                            auto client = std::make_shared<UdpClient>(recv_id,
                                [self, recv_id]() {
                                    std::cout << "received disconnect signal\n";
                                    if (self->m_clients.contains(recv_id)) {        
                                        self->m_clients.erase(recv_id);
                                    }      
                                }
                            );
                            this->m_clients.insert(std::pair<int32_t, std::shared_ptr<UdpClient>>(recv_id, client));
                            client->enqueue(buffer);

                            std::cout << "connected clients = " << this->m_clients.size() << "\n";
                        }
                    }

                    // Continue receiving...
                    this->handle_receive();
                }
            );    
        }

};



int main(int argc, char * argv[])
{    
    auto server = std::make_shared<UdpServer>(9000);
    server->listen();

    return 0;
}

编译命令

修正原命令中的无效参数:

g++ -o test main.cpp -lboost_system -lboost_date_time -lboost_thread -pthread

测试方法

nc -u 127.0.0.1 9000 < /dev/random

运行1秒后终止进程,服务器不会崩溃,5秒后会自动移除超时客户端。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 06:12:37