调用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
相关产品推荐
相关产品推荐

