22 virtual void Execute(
size_t id)
noexcept = 0;
24 template <
typename Lambda,
typename ReturnType,
typename... Args>
32 if constexpr (std::is_same_v<ReturnType, void>)
46 using JobQueues = std::array<JobQueue, kJobPriorityCount>;
71 if (!
mAccepting.load(std::memory_order_relaxed))
83 template <
typename T,
typename... Args>
84 requires std::is_base_of_v<Job, T>
88 CHECK_MSG(submitting,
"ThreadPool shutting down");
99 auto task = ConstructUniqueBase<Job, T>(allocator, std::forward<Args>(args)...);
100 mTotal.fetch_add(1, std::memory_order_release);
102 CHECK_MSG(queued,
"ThreadPool job queue full");
105 mTotal.fetch_sub(1, std::memory_order_release);
109 mWakeEpoch.fetch_add(1, std::memory_order_release);
113 template <
typename Lambda,
typename... Args>
116 auto LambdaFn = [func = std::forward<Lambda>(func), ... args = args] {
return func(args...); };
117 using LambdaType =
decltype(LambdaFn);
118 using ReturnType =
decltype(LambdaFn());
120 auto fut = job.
mPromise.get_future();
121 PushImplInternal<LambdaJob<LambdaType, ReturnType>>(priority, jobAllocator, std::move(job));
122 return std::move(fut);
127 template <
typename T,
typename... Args>
128 requires std::is_base_of_v<Job, T>
131 PushImplInternal<T>(priority,
nullptr, std::forward<Args>(args)...);
133 template <
typename T,
typename... Args>
134 requires std::is_base_of_v<Job, T>
143 template <
typename T,
typename... Args>
144 requires std::is_base_of_v<Job, T>
149 template <
typename T,
typename... Args>
150 requires std::is_base_of_v<Job, T>
153 PushImplInternal<T>(priority, jobAllocator, std::forward<Args>(args)...);
159 template <
typename Lambda,
typename... Args>
164 template <
typename Lambda,
typename... Args>
165 auto Push(Lambda&& func, Args
const&... args)
173 template <
typename Lambda,
typename... Args>
178 template <
typename Lambda,
typename... Args>
181 return PushLambdaInternal(priority, jobAllocator, std::forward<Lambda>(func), args...);
207 return mTotal.load(std::memory_order_relaxed) -
mComplete.load(std::memory_order_relaxed);
215 const static size_t CalcTaskSize(
size_t size) {
return std::bit_ceil(std::max<size_t>(size, 1)); }
#define CHECK_MSG(expr, format_str,...)
Allocator interface (noexcept)
Definition Allocator.hpp:29
Atomic, bounded multi-producer multi-consumer FIFO ring buffer with a fixed maximum size.
Definition AtomicQueue.hpp:85
Atomic, lock-free Thread Pool implementation with fixed bounds.
Definition ThreadPool.hpp:51
size_t GetTotalJobCount() const noexcept
Definition ThreadPool.hpp:210
auto Push(JobPriority priority, Lambda &&func, Args const &... args)
Push a lambda job to the thread pool.
Definition ThreadPool.hpp:160
static constexpr size_t PriorityIndex(JobPriority priority) noexcept
Definition ThreadPool.hpp:82
Atomic< bool > mAccepting
Definition ThreadPool.hpp:54
auto PushLambdaInternal(JobPriority priority, Allocator *jobAllocator, Lambda &&func, Args const &... args)
Definition ThreadPool.hpp:114
size_t GetParallelForConcurrency() const noexcept
Number of distinct worker ids a ParallelFor functor may see (workers + the participating caller)....
Definition ThreadPool.hpp:191
Allocator * mAllocator
Definition ThreadPool.hpp:52
~ThreadPool()
Join accepted jobs and stop all workers.
Definition ThreadPool.cpp:65
static const size_t CalcTaskSize(size_t size)
Definition ThreadPool.hpp:215
void PushImpl(JobPriority priority, Args &&... args)
Definition ThreadPool.hpp:129
void Shutdown()
Stop accepting work, drain accepted jobs, and stop all workers.
Definition ThreadPool.cpp:26
Atomic< size_t > mProgressEpoch
Definition ThreadPool.hpp:60
void PushImplInternal(JobPriority priority, Allocator *jobAllocator, Args &&... args)
Definition ThreadPool.hpp:85
bool BeginSubmit() noexcept
Definition ThreadPool.hpp:68
auto Push(Lambda &&func, Args const &... args)
Definition ThreadPool.hpp:165
Atomic< size_t > mComplete
Definition ThreadPool.hpp:61
void PushImpl(Args &&... args)
Definition ThreadPool.hpp:135
String mName
Definition ThreadPool.hpp:53
size_t mSubmitting
Definition ThreadPool.hpp:57
Atomic< size_t > mTotal
Definition ThreadPool.hpp:62
Atomic< bool > mShutdown
Definition ThreadPool.hpp:58
size_t GetWorkerCount() const noexcept
Number of worker threads. Worker ids passed to Execute are in [0, this).
Definition ThreadPool.hpp:185
void PushImplAlloc(Allocator *jobAllocator, Args &&... args)
Push a job with an explicit allocator for the job object.
Definition ThreadPool.hpp:145
size_t GetCompletedJobCount() const noexcept
Definition ThreadPool.hpp:209
Vector< Thread > mThreads
Definition ThreadPool.hpp:66
auto PushAlloc(JobPriority priority, Allocator *jobAllocator, Lambda &&func, Args const &... args)
Definition ThreadPool.hpp:179
void Join()
Wait for all scheduled jobs to complete.
Definition ThreadPool.cpp:56
void ThreadPoolWorker(size_t id)
Definition ThreadPool.cpp:74
void EndSubmit() noexcept
Definition ThreadPool.hpp:76
Mutex mSubmitMutex
Definition ThreadPool.hpp:55
auto PushAlloc(Allocator *jobAllocator, Lambda &&func, Args const &... args)
Push a lambda job with an explicit allocator for the job object.
Definition ThreadPool.hpp:174
CondVar mSubmitCV
Definition ThreadPool.hpp:56
void PushImplAlloc(JobPriority priority, Allocator *jobAllocator, Args &&... args)
Definition ThreadPool.hpp:151
Atomic< size_t > mWakeEpoch
Definition ThreadPool.hpp:59
size_t GetPendingJobCount() const noexcept
Definition ThreadPool.hpp:205
JobQueues mJobs
Definition ThreadPool.hpp:64
Lock-free atomic primitives and implementations of data structures.
Definition Allocator.hpp:6
std::vector< T, StlAllocator< T > > Vector
std::vector with explicit Foundation::Core::StlAllocator constructor
Definition Container.hpp:149
std::mutex Mutex
Definition Thread.hpp:10
JobPriority
Definition ThreadPool.hpp:13
std::array< JobQueue, kJobPriorityCount > JobQueues
Definition ThreadPool.hpp:46
std::condition_variable CondVar
Definition Thread.hpp:9
std::atomic< T > Atomic
Alias of std::atomic<T>.
Definition Atomic.hpp:26
std::basic_string< char, std::char_traits< char >, StlDefaultAllocator< char > > String
Alias for std::basic_string<char>, without an explicit allocator constructor.
Definition Container.hpp:120
constexpr size_t kJobPriorityCount
Definition ThreadPool.hpp:18
std::basic_string_view< char > StringView
Alias for std::basic_string_view<char>
Definition Container.hpp:56
std::promise< T > Promise
Definition Thread.hpp:6
Definition ThreadPool.hpp:20
virtual void Execute(size_t id) noexcept=0
Definition ThreadPool.hpp:26
void Execute(size_t) noexcept override
Definition ThreadPool.hpp:30
LambdaJob(Lambda &&func)
Definition ThreadPool.hpp:29
Lambda mFunc
Definition ThreadPool.hpp:27
Promise< ReturnType > mPromise
Definition ThreadPool.hpp:28