6#ifndef DEEPLIMA_SRC_INCLUDE_THREAD_POOL_H
7#define DEEPLIMA_SRC_INCLUDE_THREAD_POOL_H
16#include <condition_variable>
35 void init(
size_t num_threads)
37 assert(num_threads > 0);
41 for (
size_t i = 0; i < num_threads; i++)
63 for (
size_t i = 0 ; i <
m_workers.size(); i++)
84 throw std::runtime_error(
"All workers must be unjoinable (inactive threads) here.");
94 inline void push(
void* job)
96 std::lock_guard<std::mutex> l(
m_mutex);
108 std::unique_lock<std::mutex> l(
m_mutex);
146 P::run_one_job(
static_cast<P*
>(
this), worker_id, job);
159 throw std::runtime_error(
"Worker finished but stop flag isn't set.");
bool wait_for_new_job(void **job)
This will wait until a job is available and then job parameter will be set to this available which wi...
void init(size_t num_threads)
std::mutex m_mutex_notify
void wait_for_any_job_notification(const std::function< bool()> fn)
size_t get_num_threads() const
ThreadPool(size_t num_threads=0)
std::queue< void * > m_jobs
std::atomic< bool > m_stop
std::vector< std::thread > m_workers
std::condition_variable m_cv_notify
void thread_fn(size_t worker_id)
std::condition_variable m_cv