| 45 | * queued ones are completed. |
| 46 | */ |
| 47 | class ThreadPool |
| 48 | { |
| 49 | private: |
| 50 | std::string m_name; |
| 51 | Mutex m_mutex; |
| 52 | std::queue<std::packaged_task<void()>> m_work_queue GUARDED_BY(m_mutex); |
| 53 | std::condition_variable m_cv; |
| 54 | // Note: m_interrupt must be guarded by m_mutex, and cannot be replaced by an unguarded atomic bool. |
| 55 | // This ensures threads blocked on m_cv reliably observe the change and proceed correctly without missing signals. |
| 56 | // Ref: https://en.cppreference.com/w/cpp/thread/condition_variable |
| 57 | bool m_interrupt GUARDED_BY(m_mutex){false}; |
| 58 | std::vector<std::thread> m_workers GUARDED_BY(m_mutex); |
| 59 | |
| 60 | void WorkerThread() EXCLUSIVE_LOCKS_REQUIRED(!m_mutex) |
| 61 | { |
| 62 | WAIT_LOCK(m_mutex, wait_lock); |
| 63 | for (;;) { |
| 64 | std::packaged_task<void()> task; |
| 65 | { |
| 66 | // Wait only if needed; avoid sleeping when a new task was submitted while we were processing another one. |
| 67 | if (!m_interrupt && m_work_queue.empty()) { |
| 68 | // Block until the pool is interrupted or a task is available. |
| 69 | m_cv.wait(wait_lock, [&]() EXCLUSIVE_LOCKS_REQUIRED(m_mutex) { return m_interrupt || !m_work_queue.empty(); }); |
| 70 | } |
| 71 | |
| 72 | // If stopped and no work left, exit worker |
| 73 | if (m_interrupt && m_work_queue.empty()) { |
| 74 | return; |
| 75 | } |
| 76 | |
| 77 | task = std::move(m_work_queue.front()); |
| 78 | m_work_queue.pop(); |
| 79 | } |
| 80 | |
| 81 | { |
| 82 | // Execute the task without the lock |
| 83 | REVERSE_LOCK(wait_lock, m_mutex); |
| 84 | task(); |
| 85 | } |
| 86 | } |
| 87 | } |
| 88 | |
| 89 | public: |
| 90 | explicit ThreadPool(const std::string& name) : m_name(name) {} |
| 91 | |
| 92 | ~ThreadPool() |
| 93 | { |
| 94 | Stop(); // In case it hasn't been stopped. |
| 95 | } |
| 96 | |
| 97 | /** |
| 98 | * @brief Start worker threads. |
| 99 | * |
| 100 | * Creates and launches `num_workers` threads that begin executing tasks |
| 101 | * from the queue. If the pool is already started, throws. |
| 102 | * |
| 103 | * Must be called from a controller (non-worker) thread. |
| 104 | */ |
nothing calls this directly
no test coverage detected