// // Created by xtkuang on 2025/5/13. // #ifndef CMVR_ES_THREAD_POOL_H #define CMVR_ES_THREAD_POOL_H #include #include #include #include #include #include #include #include #include #include #pragma warning(disable:4996) #define GLOG_USE_GLOG_EXPORT #include class InterruptFlag { public: void request_stop(); [[nodiscard]] bool stop_requested() const; private: std::atomic flag_ {false}; }; // 提交任务(带 future 返回) class TaskHandle { public: using TaskFuncType = std::function; using ReturnType = void; TaskHandle(std::shared_ptr> promise) : future_(promise->get_future()) {} std::future& get_future() { return future_; } private: std::future future_; }; class ThreadPool { public: ThreadPool(size_t min_threads, size_t max_threads); ~ThreadPool(); // 提交任务(不带返回值) void submit(std::function task); TaskHandle submit_with_future(std::function f); void shutdown(); private: void worker_loop(size_t id); void monitor_loop(); boost::lockfree::queue*> task_queue_; std::atomic pending_tasks_; std::vector threads_; std::vector> flags_; std::mutex control_mutex_; std::condition_variable control_cv_; std::atomic shutdown_requested_ = false; const size_t min_threads_; const size_t max_threads_; std::thread monitor_thread_; }; #endif //CMVR_ES_THREAD_POOL_H