Eigen  5.0.1
 
Loading...
Searching...
No Matches
RunQueue.h
1// This file is part of Eigen, a lightweight C++ template library
2// for linear algebra.
3//
4// Copyright (C) 2016 Dmitry Vyukov <dvyukov@google.com>
5//
6// This Source Code Form is subject to the terms of the Mozilla
7// Public License v. 2.0. If a copy of the MPL was not distributed
8// with this file, You can obtain one at http://mozilla.org/MPL/2.0/.
9// SPDX-License-Identifier: MPL-2.0
10
11#ifndef EIGEN_THREADPOOL_RUNQUEUE_H
12#define EIGEN_THREADPOOL_RUNQUEUE_H
13
14// IWYU pragma: private
15#include "./InternalHeaderCheck.h"
16
17namespace Eigen {
18
19// RunQueue is a fixed-size, partially non-blocking deque of Work items.
20// Operations on front of the queue must be done by a single thread (owner),
21// operations on back of the queue can be done by multiple threads concurrently.
22//
23// Algorithm outline:
24// All remote threads operating on the queue back are serialized by a mutex.
25// This ensures that at most two threads access state: owner and one remote
26// thread (Size aside). The algorithm ensures that the occupied region of the
27// underlying array is logically continuous (can wraparound, but no stray
28// occupied elements). Owner operates on one end of this region, remote thread
29// operates on the other end. Synchronization between these threads
30// (potential consumption of the last element and take up of the last empty
31// element) happens by means of state variable in each element. States are:
32// empty, busy (in process of insertion of removal) and ready. Threads claim
33// elements (empty->busy and ready->busy transitions) by means of a CAS
34// operation. The finishing transition (busy->empty and busy->ready) are done
35// with plain store as the element is exclusively owned by the current thread.
36//
37// Note: we could permit only pointers as elements, then we would not need
38// separate state variable as null/non-null pointer value would serve as state,
39// but that would require malloc/free per operation for large, complex values
40// (and this is designed to store std::function<()>).
41template <typename Work, unsigned kSize>
42class RunQueue {
43 public:
44 RunQueue() : front_(0), back_(0) {
45 // require power-of-two for fast masking
46 eigen_plain_assert((kSize & (kSize - 1)) == 0);
47 eigen_plain_assert(kSize > 2); // why would you do this?
48 eigen_plain_assert(kSize <= (64 << 10)); // leave enough space for counter
49 for (unsigned i = 0; i < kSize; i++) array_[i].state.store(kEmpty, std::memory_order_relaxed);
50 }
51
52 ~RunQueue() { eigen_plain_assert(Size() == 0); }
53
54 // PushFront inserts w at the beginning of the queue.
55 // If queue is full returns w, otherwise returns default-constructed Work.
56 Work PushFront(Work w) {
57 unsigned front = front_.load(std::memory_order_relaxed);
58 Elem* e = &array_[front & kMask];
59 uint8_t s = e->state.load(std::memory_order_relaxed);
60 if (s != kEmpty || !e->state.compare_exchange_strong(s, kBusy, std::memory_order_acquire)) return w;
61 front_.store(front + 1 + (kSize << 1), std::memory_order_release);
62 e->w = std::move(w);
63 e->state.store(kReady, std::memory_order_release);
64 return Work();
65 }
66
67 // PopFront removes and returns the first element in the queue.
68 // If the queue was empty returns default-constructed Work.
69 Work PopFront() {
70 unsigned front = front_.load(std::memory_order_relaxed);
71 Elem* e = &array_[(front - 1) & kMask];
72 uint8_t s = e->state.load(std::memory_order_relaxed);
73 if (s != kReady || !e->state.compare_exchange_strong(s, kBusy, std::memory_order_acquire)) return Work();
74 Work w = std::move(e->w);
75 e->state.store(kEmpty, std::memory_order_release);
76 front = ((front - 1) & kMask2) | (front & ~kMask2);
77 front_.store(front, std::memory_order_release);
78 return w;
79 }
80
81 // PushBack adds w at the end of the queue.
82 // If queue is full returns w, otherwise returns default-constructed Work.
83 Work PushBack(Work w) {
84 EIGEN_MUTEX_LOCK lock(mutex_);
85 unsigned back = back_.load(std::memory_order_relaxed);
86 Elem* e = &array_[(back - 1) & kMask];
87 uint8_t s = e->state.load(std::memory_order_relaxed);
88 if (s != kEmpty || !e->state.compare_exchange_strong(s, kBusy, std::memory_order_acquire)) return w;
89 back = ((back - 1) & kMask2) | (back & ~kMask2);
90 back_.store(back, std::memory_order_release);
91 e->w = std::move(w);
92 e->state.store(kReady, std::memory_order_release);
93 return Work();
94 }
95
96 // PopBack removes and returns the last element in the queue.
97 Work PopBack() {
98 if (Empty()) return Work();
99 EIGEN_MUTEX_LOCK lock(mutex_);
100 unsigned back = back_.load(std::memory_order_relaxed);
101 Elem* e = &array_[back & kMask];
102 uint8_t s = e->state.load(std::memory_order_relaxed);
103 if (s != kReady || !e->state.compare_exchange_strong(s, kBusy, std::memory_order_acquire)) return Work();
104 Work w = std::move(e->w);
105 e->state.store(kEmpty, std::memory_order_release);
106 back_.store(back + 1 + (kSize << 1), std::memory_order_release);
107 return w;
108 }
109
110 // PopBackHalf removes and returns half last elements in the queue.
111 // Returns number of elements removed.
112 unsigned PopBackHalf(std::vector<Work>* result) {
113 if (Empty()) return 0;
114 EIGEN_MUTEX_LOCK lock(mutex_);
115 unsigned back = back_.load(std::memory_order_relaxed);
116 unsigned size = Size();
117 unsigned mid = back;
118 if (size > 1) mid = back + (size - 1) / 2;
119 unsigned n = 0;
120 unsigned start = 0;
121 for (; static_cast<int>(mid - back) >= 0; mid--) {
122 Elem* e = &array_[mid & kMask];
123 uint8_t s = e->state.load(std::memory_order_relaxed);
124 if (n == 0) {
125 if (s != kReady || !e->state.compare_exchange_strong(s, kBusy, std::memory_order_acquire)) continue;
126 start = mid;
127 } else {
128 // Note: no need to store temporal kBusy, we exclusively own these
129 // elements.
130 eigen_plain_assert(s == kReady);
131 }
132 result->push_back(std::move(e->w));
133 e->state.store(kEmpty, std::memory_order_release);
134 n++;
135 }
136 if (n != 0) back_.store(start + 1 + (kSize << 1), std::memory_order_release);
137 return n;
138 }
139
140 // Size returns current queue size.
141 // Can be called by any thread at any time.
142 unsigned Size() const { return SizeOrNotEmpty<true>(); }
143
144 // Empty tests whether container is empty.
145 // Can be called by any thread at any time.
146 bool Empty() const { return SizeOrNotEmpty<false>() == 0; }
147
148 // Delete all the elements from the queue.
149 void Flush() {
150 while (!Empty()) {
151 PopFront();
152 }
153 }
154
155 private:
156 static const unsigned kMask = kSize - 1;
157 static const unsigned kMask2 = (kSize << 1) - 1;
158
159 enum State {
160 kEmpty,
161 kBusy,
162 kReady,
163 };
164
165 struct Elem {
166 std::atomic<uint8_t> state;
167 Work w;
168 };
169
170 // Low log(kSize) + 1 bits in front_ and back_ contain rolling index of
171 // front/back, respectively. The remaining bits contain modification counters
172 // that are incremented on Push operations. This allows us to (1) distinguish
173 // between empty and full conditions (if we would use log(kSize) bits for
174 // position, these conditions would be indistinguishable); (2) obtain
175 // consistent snapshot of front_/back_ for Size operation using the
176 // modification counters.
177 EIGEN_ALIGN_TO_AVOID_FALSE_SHARING std::atomic<unsigned> front_;
178 EIGEN_ALIGN_TO_AVOID_FALSE_SHARING std::atomic<unsigned> back_;
179 EIGEN_MUTEX mutex_; // guards `PushBack` and `PopBack` (accesses `back_`)
180
181 EIGEN_ALIGN_TO_AVOID_FALSE_SHARING Elem array_[kSize];
182
183 // SizeOrNotEmpty returns current queue size; if NeedSizeEstimate is false,
184 // only whether the size is 0 is guaranteed to be correct.
185 // Can be called by any thread at any time.
186 template <bool NeedSizeEstimate>
187 unsigned SizeOrNotEmpty() const {
188 // Emptiness plays critical role in thread pool blocking. So we go to great
189 // effort to not produce false positives (claim non-empty queue as empty).
190 unsigned front = front_.load(std::memory_order_acquire);
191 for (;;) {
192 // Capture a consistent snapshot of front/back.
193 unsigned back = back_.load(std::memory_order_acquire);
194 unsigned front1 = front_.load(std::memory_order_relaxed);
195 if (front != front1) {
196 front = front1;
197 std::atomic_thread_fence(std::memory_order_acquire);
198 continue;
199 }
200 EIGEN_IF_CONSTEXPR (NeedSizeEstimate) {
201 return CalculateSize(front, back);
202 } else {
203 // This value will be 0 if the queue is empty, and undefined otherwise.
204 unsigned maybe_zero = ((front ^ back) & kMask2);
205 // Queue size estimate must agree with maybe zero check on the queue
206 // empty/non-empty state.
207 eigen_assert((CalculateSize(front, back) == 0) == (maybe_zero == 0));
208 return maybe_zero;
209 }
210 }
211 }
212
213 EIGEN_ALWAYS_INLINE unsigned CalculateSize(unsigned front, unsigned back) const {
214 int size = (front & kMask2) - (back & kMask2);
215 // Fix overflow.
216 if (EIGEN_PREDICT_FALSE(size < 0)) size += 2 * kSize;
217 // Order of modification in push/pop is crafted to make the queue look
218 // larger than it is during concurrent modifications. E.g. push can
219 // increment size before the corresponding pop has decremented it.
220 // So the computed size can be up to kSize + 1, fix it.
221 if (EIGEN_PREDICT_FALSE(size > static_cast<int>(kSize))) size = kSize;
222 return static_cast<unsigned>(size);
223 }
224
225 RunQueue(const RunQueue&) = delete;
226 void operator=(const RunQueue&) = delete;
227};
228
229} // namespace Eigen
230
231#endif // EIGEN_THREADPOOL_RUNQUEUE_H