Foundation
Loading...
Searching...
No Matches
Public Member Functions | Static Public Member Functions | Private Member Functions | Static Private Member Functions | Private Attributes | List of all members
Foundation::Core::ThreadPool Class Reference

Atomic, lock-free Thread Pool implementation with fixed bounds. More...

#include <ThreadPool.hpp>

Public Member Functions

 ThreadPool (size_t numThreads, size_t maxTasks, Allocator *alloc, StringView name="ThreadPool")
 
template<typename T , typename... Args>
requires std::is_base_of_v<Job, T>
void PushImpl (JobPriority priority, Args &&... args)
 
template<typename T , typename... Args>
requires std::is_base_of_v<Job, T>
void PushImpl (Args &&... args)
 
template<typename T , typename... Args>
requires std::is_base_of_v<Job, T>
void PushImplAlloc (Allocator *jobAllocator, Args &&... args)
 Push a job with an explicit allocator for the job object.
 
template<typename T , typename... Args>
requires std::is_base_of_v<Job, T>
void PushImplAlloc (JobPriority priority, Allocator *jobAllocator, Args &&... args)
 
template<typename Lambda , typename... Args>
auto Push (JobPriority priority, Lambda &&func, Args const &... args)
 Push a lambda job to the thread pool.
 
template<typename Lambda , typename... Args>
auto Push (Lambda &&func, Args const &... args)
 
template<typename Lambda , typename... Args>
auto PushAlloc (Allocator *jobAllocator, Lambda &&func, Args const &... args)
 Push a lambda job with an explicit allocator for the job object.
 
template<typename Lambda , typename... Args>
auto PushAlloc (JobPriority priority, Allocator *jobAllocator, Lambda &&func, Args const &... args)
 
size_t GetWorkerCount () const noexcept
 Number of worker threads. Worker ids passed to Execute are in [0, this).
 
size_t GetParallelForConcurrency () const noexcept
 Number of distinct worker ids a ParallelFor functor may see (workers + the participating caller). Size per-worker scratch to this.
 
void Shutdown ()
 Stop accepting work, drain accepted jobs, and stop all workers.
 
void Join ()
 Wait for all scheduled jobs to complete.
 
 ~ThreadPool ()
 Join accepted jobs and stop all workers.
 
size_t GetPendingJobCount () const noexcept
 
size_t GetCompletedJobCount () const noexcept
 
size_t GetTotalJobCount () const noexcept
 

Static Public Member Functions

static const size_t CalcTaskSize (size_t size)
 

Private Member Functions

void ThreadPoolWorker (size_t id)
 
bool BeginSubmit () noexcept
 
void EndSubmit () noexcept
 
template<typename T , typename... Args>
requires std::is_base_of_v<Job, T>
void PushImplInternal (JobPriority priority, Allocator *jobAllocator, Args &&... args)
 
template<typename Lambda , typename... Args>
auto PushLambdaInternal (JobPriority priority, Allocator *jobAllocator, Lambda &&func, Args const &... args)
 

Static Private Member Functions

static constexpr size_t PriorityIndex (JobPriority priority) noexcept
 

Private Attributes

AllocatormAllocator
 
String mName
 
Atomic< bool > mAccepting {true}
 
Mutex mSubmitMutex
 
CondVar mSubmitCV
 
size_t mSubmitting {}
 
Atomic< bool > mShutdown {false}
 
Atomic< size_t > mWakeEpoch {0}
 
Atomic< size_t > mProgressEpoch {0}
 
Atomic< size_t > mComplete {0}
 
Atomic< size_t > mTotal {0}
 
JobQueues mJobs
 
Vector< ThreadmThreads
 

Detailed Description

Atomic, lock-free Thread Pool implementation with fixed bounds.

Constructor & Destructor Documentation

◆ ThreadPool()

Foundation::Core::ThreadPool::ThreadPool ( size_t  numThreads,
size_t  maxTasks,
Allocator alloc,
StringView  name = "ThreadPool" 
)

◆ ~ThreadPool()

Foundation::Core::ThreadPool::~ThreadPool ( )

Join accepted jobs and stop all workers.

Member Function Documentation

◆ BeginSubmit()

bool Foundation::Core::ThreadPool::BeginSubmit ( )
inlineprivatenoexcept

◆ CalcTaskSize()

static const size_t Foundation::Core::ThreadPool::CalcTaskSize ( size_t  size)
inlinestatic

Aligns a number to upper, closest power of 2 so that it's a valid maxTasks size.

◆ EndSubmit()

void Foundation::Core::ThreadPool::EndSubmit ( )
inlineprivatenoexcept

◆ GetCompletedJobCount()

size_t Foundation::Core::ThreadPool::GetCompletedJobCount ( ) const
inlinenoexcept

◆ GetParallelForConcurrency()

size_t Foundation::Core::ThreadPool::GetParallelForConcurrency ( ) const
inlinenoexcept

Number of distinct worker ids a ParallelFor functor may see (workers + the participating caller). Size per-worker scratch to this.

◆ GetPendingJobCount()

size_t Foundation::Core::ThreadPool::GetPendingJobCount ( ) const
inlinenoexcept

◆ GetTotalJobCount()

size_t Foundation::Core::ThreadPool::GetTotalJobCount ( ) const
inlinenoexcept

◆ GetWorkerCount()

size_t Foundation::Core::ThreadPool::GetWorkerCount ( ) const
inlinenoexcept

Number of worker threads. Worker ids passed to Execute are in [0, this).

◆ Join()

void Foundation::Core::ThreadPool::Join ( )

Wait for all scheduled jobs to complete.

Note
Concurrent submission must be externally synchronized with this call.

◆ PriorityIndex()

static constexpr size_t Foundation::Core::ThreadPool::PriorityIndex ( JobPriority  priority)
inlinestaticconstexprprivatenoexcept

◆ Push() [1/2]

template<typename Lambda , typename... Args>
auto Foundation::Core::ThreadPool::Push ( JobPriority  priority,
Lambda &&  func,
Args const &...  args 
)
inline

Push a lambda job to the thread pool.

Returns
Future<func ReturnType> that will be set when the job is completed.

◆ Push() [2/2]

template<typename Lambda , typename... Args>
auto Foundation::Core::ThreadPool::Push ( Lambda &&  func,
Args const &...  args 
)
inline

◆ PushAlloc() [1/2]

template<typename Lambda , typename... Args>
auto Foundation::Core::ThreadPool::PushAlloc ( Allocator jobAllocator,
Lambda &&  func,
Args const &...  args 
)
inline

Push a lambda job with an explicit allocator for the job object.

Parameters
jobAllocatorOptional allocator for the job object. If null, the thread pool allocator is used.

◆ PushAlloc() [2/2]

template<typename Lambda , typename... Args>
auto Foundation::Core::ThreadPool::PushAlloc ( JobPriority  priority,
Allocator jobAllocator,
Lambda &&  func,
Args const &...  args 
)
inline

◆ PushImpl() [1/2]

template<typename T , typename... Args>
requires std::is_base_of_v<Job, T>
void Foundation::Core::ThreadPool::PushImpl ( Args &&...  args)
inline

◆ PushImpl() [2/2]

template<typename T , typename... Args>
requires std::is_base_of_v<Job, T>
void Foundation::Core::ThreadPool::PushImpl ( JobPriority  priority,
Args &&...  args 
)
inline

◆ PushImplAlloc() [1/2]

template<typename T , typename... Args>
requires std::is_base_of_v<Job, T>
void Foundation::Core::ThreadPool::PushImplAlloc ( Allocator jobAllocator,
Args &&...  args 
)
inline

Push a job with an explicit allocator for the job object.

Parameters
jobAllocatorOptional allocator for the job object. If null, the thread pool allocator is used.

◆ PushImplAlloc() [2/2]

template<typename T , typename... Args>
requires std::is_base_of_v<Job, T>
void Foundation::Core::ThreadPool::PushImplAlloc ( JobPriority  priority,
Allocator jobAllocator,
Args &&...  args 
)
inline

◆ PushImplInternal()

template<typename T , typename... Args>
requires std::is_base_of_v<Job, T>
void Foundation::Core::ThreadPool::PushImplInternal ( JobPriority  priority,
Allocator jobAllocator,
Args &&...  args 
)
inlineprivate

◆ PushLambdaInternal()

template<typename Lambda , typename... Args>
auto Foundation::Core::ThreadPool::PushLambdaInternal ( JobPriority  priority,
Allocator jobAllocator,
Lambda &&  func,
Args const &...  args 
)
inlineprivate

◆ Shutdown()

void Foundation::Core::ThreadPool::Shutdown ( )

Stop accepting work, drain accepted jobs, and stop all workers.

◆ ThreadPoolWorker()

void Foundation::Core::ThreadPool::ThreadPoolWorker ( size_t  id)
private

Member Data Documentation

◆ mAccepting

Atomic<bool> Foundation::Core::ThreadPool::mAccepting {true}
private

◆ mAllocator

Allocator* Foundation::Core::ThreadPool::mAllocator
private

◆ mComplete

Atomic<size_t> Foundation::Core::ThreadPool::mComplete {0}
private

◆ mJobs

JobQueues Foundation::Core::ThreadPool::mJobs
private

◆ mName

String Foundation::Core::ThreadPool::mName
private

◆ mProgressEpoch

Atomic<size_t> Foundation::Core::ThreadPool::mProgressEpoch {0}
private

◆ mShutdown

Atomic<bool> Foundation::Core::ThreadPool::mShutdown {false}
private

◆ mSubmitCV

CondVar Foundation::Core::ThreadPool::mSubmitCV
private

◆ mSubmitMutex

Mutex Foundation::Core::ThreadPool::mSubmitMutex
private

◆ mSubmitting

size_t Foundation::Core::ThreadPool::mSubmitting {}
private

◆ mThreads

Vector<Thread> Foundation::Core::ThreadPool::mThreads
private

◆ mTotal

Atomic<size_t> Foundation::Core::ThreadPool::mTotal {0}
private

◆ mWakeEpoch

Atomic<size_t> Foundation::Core::ThreadPool::mWakeEpoch {0}
private

The documentation for this class was generated from the following files: