5#include <condition_variable>
14namespace Vkm::Engine {
42 Batch(
size_t count,
const Fn& task)
43 : m_run([](const void* body, size_t index) { (*
static_cast<const Fn*
>(body))(index); })
50 Batch(
size_t count,
const Fn && task) =
delete;
61 friend class ThreadPool;
64 void (*m_run)(
const void* task,
size_t index);
68 std::exception_ptr m_error;
70 std::atomic<size_t> m_pending{0};
86 static ThreadPool&
get();
95 void addTask(std::function<
void()> && task);
135 size_t threadCount()
const {
return m_threads.size(); }
144 std::function<void()> function;
145 Batch* batch =
nullptr;
150 ThreadPool(
size_t threadCount);
166 void runIndex(Batch& batch,
size_t index);
176 void retire(std::atomic<size_t>& pending);
179 std::atomic<bool> m_running;
181 std::vector<std::thread> m_threads;
183 std::deque<QueuedTask> m_frameTasks;
184 std::deque<QueuedTask> m_backgroundTasks;
186 std::mutex m_tasksMutex;
187 std::condition_variable m_tasksCV;
188 std::condition_variable m_doneCV;
203template<
class Function>
204void parallelFor(
size_t count,
size_t grain, Function&& function) {
209 grain = std::max(grain,
size_t(1));
211 auto invokeAt = [&](
size_t i) {
212 if constexpr (std::is_invocable_v<Function, size_t>) {
220 for (
size_t i = 0; i < count; ++i) invokeAt(i);
227 if (pool.threadCount() == 0) {
228 for (
size_t i = 0; i < count; ++i) invokeAt(i);
233 const size_t chunks = (count - 1) / grain + 1;
234 const auto runChunk = [&](
size_t chunk) {
235 const size_t end = std::min(count, (chunk + 1) * grain);
236 for (
size_t index = chunk * grain; index < end; ++index) invokeAt(index);
238 const auto runQueued = [&](
size_t task) { runChunk(task + 1); };
241 pool.addBatch(batch);
249 pool.waitForBatch(batch);
254 pool.waitForBatch(batch);
266template<
class Function>
267void parallelFor(
size_t count, Function&& function) {
272 constexpr size_t MIN_PARALLEL = 2048;
275 const size_t grain = (count < MIN_PARALLEL)
277 : count / (pool.threadCount() + 1);
279 parallelFor(count, grain, function);
A run of tasks addressed by index, queued as one entry.
Definition thread_pool.h:32
Batch(size_t count, const Fn &task)
A batch of count tasks, each a call of task with its index.
Definition thread_pool.h:42
Batch(size_t count, const Fn &&task)=delete
Refused: a temporary task is gone before a worker reaches it.
Fixed-size pool of worker threads draining two task queues.
Definition thread_pool.h:24
static bool isWorkerThread()
True when called from a thread owned by the pool.
void shutdown()
Join the workers now, ahead of the pool's own destruction.
void addTask(std::function< void()> &&task)
Enqueue a single task and wake one worker.
void waitForBatch(Batch &batch)
Block the caller until every index of batch has retired.
static ThreadPool & get()
Access the process-wide thread pool, constructed on first use.
void addBatch(Batch &batch)
Queue every index of batch as one entry and wake all workers.