mirror of
https://github.com/odin-lang/Odin.git
synced 2026-10-09 14:22:18 -04:00
Thread pool: wake one sleeping worker per task instead of all, spin briefly before sleeping; walk the procedure dependency tree in chunks
This commit is contained in:
1 parent
c494c7a503
commit
0de3f92449
2 files changed
+92
-65
No files matched your search
+68
-50
@@ -29,20 +29,20 @@ enum GrabState {
|
||||
Grab_Failed = 2,
|
||||
};
|
||||
|
||||
enum BroadcastWaitState {
|
||||
Nobody_Waiting = 0,
|
||||
Someone_Waiting = 1,
|
||||
};
|
||||
|
||||
struct ThreadPool {
|
||||
gbAllocator threads_allocator;
|
||||
Slice<Thread> threads;
|
||||
std::atomic<bool> running;
|
||||
|
||||
Futex tasks_available;
|
||||
Futex tasks_left;
|
||||
Futex tasks_available; // bumped to wake a sleeping worker
|
||||
std::atomic<i32> sleeping; // workers asleep on `tasks_available`, or about to be
|
||||
Futex tasks_left;
|
||||
};
|
||||
|
||||
// NOTE(bill): how many times an idle worker looks for a task before sleeping, so one adding small tasks
|
||||
// one after another keeps the workers busy rather than waking one for each
|
||||
enum { THREAD_POOL_SPIN_COUNT = 64 };
|
||||
|
||||
gb_internal isize current_thread_index(void) {
|
||||
return current_thread ? current_thread->idx : 0;
|
||||
}
|
||||
@@ -69,8 +69,8 @@ gb_internal void thread_pool_destroy(ThreadPool *pool) {
|
||||
|
||||
for_array_off(i, 1, pool->threads) {
|
||||
Thread *t = &pool->threads[i];
|
||||
pool->tasks_available.store(Nobody_Waiting);
|
||||
futex_broadcast(&t->pool->tasks_available);
|
||||
pool->tasks_available.fetch_add(1);
|
||||
futex_broadcast(&pool->tasks_available);
|
||||
thread_join_and_destroy(t);
|
||||
}
|
||||
|
||||
@@ -103,9 +103,13 @@ void thread_pool_queue_push(Thread *thread, WorkerTask task) {
|
||||
thread->queue.bottom.store(bot + 1, std::memory_order_relaxed);
|
||||
|
||||
thread->pool->tasks_left.fetch_add(1, std::memory_order_release);
|
||||
i32 state = Someone_Waiting;
|
||||
if (thread->pool->tasks_available.compare_exchange_strong(state, Nobody_Waiting)) {
|
||||
futex_broadcast(&thread->pool->tasks_available);
|
||||
|
||||
// NOTE(bill): one sleeping worker per task; waking them all made a loop adding small tasks mostly
|
||||
// wake-ups, as they were back asleep before the next. The fence pairs with the one in the worker loop.
|
||||
std::atomic_thread_fence(std::memory_order_seq_cst);
|
||||
if (thread->pool->sleeping.load(std::memory_order_relaxed) > 0) {
|
||||
thread->pool->tasks_available.fetch_add(1);
|
||||
futex_signal(&thread->pool->tasks_available);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -165,6 +169,40 @@ GrabState thread_pool_queue_steal(Thread *thread, WorkerTask *task) {
|
||||
return ret;
|
||||
}
|
||||
|
||||
gb_internal bool thread_pool_queue_has_tasks(Thread *thread) {
|
||||
return thread->queue.top.load(std::memory_order_acquire) < thread->queue.bottom.load(std::memory_order_acquire);
|
||||
}
|
||||
|
||||
// Runs a task from another thread's queue; false if none had one
|
||||
gb_internal bool thread_pool_steal(ThreadPool *pool) {
|
||||
usize idx = cast(usize)current_thread->idx;
|
||||
for_array(i, pool->threads) {
|
||||
idx = (idx + 1) % cast(usize)pool->threads.count;
|
||||
Thread *thread = &pool->threads.data[idx];
|
||||
if (!thread_pool_queue_has_tasks(thread)) {
|
||||
continue;
|
||||
}
|
||||
|
||||
WorkerTask task;
|
||||
switch (thread_pool_queue_steal(thread, &task)) {
|
||||
case Grab_Empty:
|
||||
continue;
|
||||
case Grab_Success:
|
||||
task.do_work(task.data);
|
||||
pool->tasks_left.fetch_sub(1, std::memory_order_release);
|
||||
|
||||
if (pool->tasks_left.load(std::memory_order_acquire) == 0) {
|
||||
futex_signal(&pool->tasks_left);
|
||||
}
|
||||
return true;
|
||||
case Grab_Failed:
|
||||
// NOTE: another thread took it, so there may be more
|
||||
return true;
|
||||
}
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
gb_internal bool thread_pool_add_task(ThreadPool *pool, WorkerTaskProc *proc, void *data) {
|
||||
WorkerTask task = {};
|
||||
task.do_work = proc;
|
||||
@@ -247,7 +285,6 @@ gb_internal THREAD_PROC(thread_pool_thread_proc) {
|
||||
while (pool->running.load(std::memory_order_seq_cst)) {
|
||||
// If we've got tasks to process, work through them
|
||||
usize finished_tasks = 0;
|
||||
i32 state;
|
||||
|
||||
while (!thread_pool_queue_take(current_thread, &task)) {
|
||||
task.do_work(task.data);
|
||||
@@ -260,48 +297,29 @@ gb_internal THREAD_PROC(thread_pool_thread_proc) {
|
||||
}
|
||||
|
||||
// If there's still work somewhere and we don't have it, steal it
|
||||
if (pool->tasks_left.load(std::memory_order_acquire)) {
|
||||
usize idx = cast(usize)current_thread->idx;
|
||||
for_array(i, pool->threads) {
|
||||
if (pool->tasks_left.load(std::memory_order_acquire) == 0) {
|
||||
break;
|
||||
}
|
||||
|
||||
idx = (idx + 1) % cast(usize)pool->threads.count;
|
||||
|
||||
Thread *thread = &pool->threads.data[idx];
|
||||
WorkerTask task;
|
||||
|
||||
GrabState ret = thread_pool_queue_steal(thread, &task);
|
||||
switch (ret) {
|
||||
case Grab_Empty:
|
||||
continue;
|
||||
case Grab_Success:
|
||||
task.do_work(task.data);
|
||||
pool->tasks_left.fetch_sub(1, std::memory_order_release);
|
||||
|
||||
if (pool->tasks_left.load(std::memory_order_acquire) == 0) {
|
||||
futex_signal(&pool->tasks_left);
|
||||
}
|
||||
|
||||
/*fallthrough*/
|
||||
case Grab_Failed:
|
||||
goto main_loop_continue;
|
||||
}
|
||||
for (isize spin = 0; spin < THREAD_POOL_SPIN_COUNT; spin++) {
|
||||
if (thread_pool_steal(pool)) {
|
||||
goto main_loop_continue;
|
||||
}
|
||||
yield_thread();
|
||||
}
|
||||
|
||||
// if we've done all our work, and there's nothing to steal, go to sleep
|
||||
pool->tasks_available.store(Someone_Waiting);
|
||||
if (!pool->running) {
|
||||
// do not leave the word published on the way out: a worker still on its way to
|
||||
// futex_wait would sleep on it, and the destroyer does not broadcast again once it
|
||||
// has committed to its first join
|
||||
pool->tasks_available.store(Nobody_Waiting);
|
||||
futex_broadcast(&pool->tasks_available);
|
||||
break;
|
||||
{
|
||||
Footex epoch = pool->tasks_available.load();
|
||||
pool->sleeping.fetch_add(1, std::memory_order_relaxed);
|
||||
std::atomic_thread_fence(std::memory_order_seq_cst);
|
||||
|
||||
// NOTE: a task added before `sleeping` was raised woke nobody, so look again
|
||||
bool has_tasks = false;
|
||||
for (Thread &t : pool->threads) {
|
||||
has_tasks |= thread_pool_queue_has_tasks(&t);
|
||||
}
|
||||
if (!has_tasks && pool->running.load()) {
|
||||
futex_wait(&pool->tasks_available, epoch);
|
||||
}
|
||||
pool->sleeping.fetch_sub(1, std::memory_order_relaxed);
|
||||
}
|
||||
futex_wait(&pool->tasks_available, Someone_Waiting);
|
||||
|
||||
main_loop_continue:;
|
||||
}
|
||||
|
||||
Reference in new issue
Block a user