64 using Metadata =
typename Policy::Metadata;
68 std::function<void()> run;
71 using Container =
typename Policy::template Container<WorkItem>;
78 std::vector<std::unique_ptr<WorkerQueue>> queues_;
79 std::atomic<size_t> queuedTasks_{0};
80 std::vector<std::thread> workers_;
81 mutable std::mutex waitMutex_;
82 std::condition_variable condition_;
83 std::atomic<bool> stop_{
false};
84 std::atomic<int> activeTasks_{0};
85 inline static thread_local Scheduler* currentScheduler_ =
nullptr;
86 inline static thread_local int workerIndex_ = -1;
87 std::atomic<size_t> nextWorker_{0};
95 void worker_thread(
size_t index) {
96 currentScheduler_ =
this;
97 workerIndex_ =
static_cast<int>(index);
100 if (!tryPopTask(index, item)) {
101 std::unique_lock<std::mutex> lock(waitMutex_);
102 condition_.wait(lock, [
this] {
103 return stop_.load(std::memory_order_acquire) ||
104 queuedTasks_.load(std::memory_order_acquire) > 0;
106 if (stop_.load(std::memory_order_acquire) &&
107 queuedTasks_.load(std::memory_order_acquire) == 0) {
108 currentScheduler_ =
nullptr;
122 if (activeTasks_.fetch_sub(1, std::memory_order_release) == 1) {
123 std::lock_guard<std::mutex> lock(waitMutex_);
124 condition_.notify_all();
130 bool tryPopTask(
size_t index, WorkItem& item) {
132 if (tryPopLocal(index, item))
return true;
133 if (trySteal(index, item))
return true;
138 bool tryPopLocal(
size_t index, WorkItem& item) {
139 WorkerQueue& queue = *queues_[index];
140 std::lock_guard<std::mutex> lock(queue.mutex);
141 if (!Policy::template popLocal<WorkItem>(queue.queue, item))
return false;
142 queuedTasks_.fetch_sub(1, std::memory_order_release);
147 bool trySteal(
size_t index, WorkItem& item) {
148 const size_t workerCount = queues_.size();
149 if (workerCount <= 1)
return false;
151 for (
size_t offset = 1; offset < workerCount; ++offset) {
152 size_t target = (index + offset) % workerCount;
153 WorkerQueue& queue = *queues_[target];
154 std::unique_lock<std::mutex> lock(queue.mutex, std::try_to_lock);
155 if (!lock.owns_lock())
continue;
156 if (!Policy::template popSteal<WorkItem>(queue.queue, item))
continue;
157 queuedTasks_.fetch_sub(1, std::memory_order_release);
164 bool isWorkerThread()
const {
return currentScheduler_ ==
this; }
166 bool enqueueImpl(Metadata metadata, std::function<
void()> work) {
167 if (stop_.load(std::memory_order_acquire))
return false;
168 WorkerQueue* targetQueue =
nullptr;
169 if (isWorkerThread()) {
170 targetQueue = queues_[
static_cast<size_t>(workerIndex_)].get();
172 const size_t target =
173 nextWorker_.fetch_add(1, std::memory_order_relaxed) % queues_.size();
174 targetQueue = queues_[target].get();
178 std::lock_guard<std::mutex> lock(targetQueue->mutex);
179 Policy::template push<WorkItem>(
180 targetQueue->queue, WorkItem{std::move(metadata), std::move(work)});
183 queuedTasks_.fetch_add(1, std::memory_order_release);
184 activeTasks_.fetch_add(1, std::memory_order_release);
186 std::lock_guard<std::mutex> lock(waitMutex_);
187 condition_.notify_one();
192 template <
typename T = Y,
typename = std::enable_if_t<std::is_
void_v<T>>>
193 void scheduleOrRunInlineImpl(Metadata metadata, std::function<
void()> job) {
194 if (stop_.load(std::memory_order_relaxed))
return;
195 if (!isWorkerThread()) {
196 enqueueImpl(std::move(metadata), std::move(job));
200 activeTasks_.fetch_add(1, std::memory_order_release);
205 if (activeTasks_.fetch_sub(1, std::memory_order_release) == 1) {
206 std::lock_guard<std::mutex> lock(waitMutex_);
207 condition_.notify_all();
218 explicit Scheduler(
size_t numThreads = std::thread::hardware_concurrency()) {
219 if (numThreads == 0) numThreads = 1;
220 queues_.reserve(numThreads);
221 for (
size_t i = 0; i < numThreads; ++i) {
222 queues_.push_back(std::make_unique<WorkerQueue>());
224 for (
size_t i = 0; i < numThreads; ++i) {
225 workers_.emplace_back(&Scheduler::worker_thread,
this, i);
236 stop_.store(
true, std::memory_order_release);
237 condition_.notify_all();
238 for (std::thread& worker : workers_) {
239 if (worker.joinable()) worker.join();
251 template <
typename P = Policy,
typename = std::enable_if_t<std::is_same_v<
252 typename P::Metadata, std::monostate>>>
254 if (stop_.load(std::memory_order_acquire)) {
255 std::promise<Y> err_promise;
256 err_promise.set_exception(std::make_exception_ptr(
257 std::runtime_error(
"Scheduler is stopping or stopped.")));
258 return err_promise.get_future();
261 auto promise = std::make_shared<std::promise<Y>>();
262 std::future<Y> future = promise->get_future();
264 auto work = [promise, job = std::move(job)]()
mutable {
266 if constexpr (std::is_void_v<Y>) {
268 promise->set_value();
270 promise->set_value(job());
274 promise->set_exception(std::current_exception());
280 if (!enqueueImpl(Metadata{}, std::move(work))) {
282 promise->set_exception(std::make_exception_ptr(
283 std::runtime_error(
"Scheduler is stopping or stopped.")));
299 template <
typename P = Policy,
typename = std::enable_if_t<!std::is_same_v<
300 typename P::Metadata, std::monostate>>>
301 std::future<Y>
schedule(Metadata metadata, std::function<Y()> job) {
302 if (stop_.load(std::memory_order_acquire)) {
303 std::promise<Y> err_promise;
304 err_promise.set_exception(std::make_exception_ptr(
305 std::runtime_error(
"Scheduler is stopping or stopped.")));
306 return err_promise.get_future();
309 auto promise = std::make_shared<std::promise<Y>>();
310 std::future<Y> future = promise->get_future();
312 auto work = [promise, job = std::move(job)]()
mutable {
314 if constexpr (std::is_void_v<Y>) {
316 promise->set_value();
318 promise->set_value(job());
322 promise->set_exception(std::current_exception());
328 if (!enqueueImpl(std::move(metadata), std::move(work))) {
330 promise->set_exception(std::make_exception_ptr(
331 std::runtime_error(
"Scheduler is stopping or stopped.")));
343 template <
typename T = Y,
344 typename = std::enable_if_t<
345 std::is_void_v<T> && std::is_same_v<Metadata, std::monostate>>>
347 scheduleOrRunInlineImpl(Metadata{}, std::move(job));
355 template <
typename T = Y,
356 typename = std::enable_if_t<
357 std::is_void_v<T> && !std::is_same_v<Metadata, std::monostate>>>
359 scheduleOrRunInlineImpl(std::move(metadata), std::move(job));
367 template <
typename T = Y,
368 typename = std::enable_if_t<
369 std::is_void_v<T> && std::is_same_v<Metadata, std::monostate>>>
371 enqueueImpl(Metadata{}, std::move(job));
379 template <
typename T = Y,
380 typename = std::enable_if_t<
381 std::is_void_v<T> && std::is_same_v<Metadata, std::monostate>>>
383 scheduleOrRunInlineImpl(Metadata{}, std::move(job));
392 std::unique_lock<std::mutex> lock(waitMutex_);
393 condition_.wait(lock, [
this] {
394 return stop_.load(std::memory_order_acquire) ||
395 (activeTasks_.load(std::memory_order_acquire) == 0 &&
396 queuedTasks_.load(std::memory_order_acquire) == 0);