Cpp-Taskflow  2.3.0
wsq.hpp
1 // 2019/05/15 - created by Tsung-Wei Huang
2 // - isolated from the original workstealing executor
3 
4 #pragma once
5 
6 #include <atomic>
7 #include <vector>
8 #include <optional>
9 
10 namespace tf {
11 
26 template <typename T>
28 
29  //constexpr static int64_t cacheline_size = 64;
30 
31  //using storage_type = std::aligned_storage_t<sizeof(T), cacheline_size>;
32 
33  struct Array {
34 
35  int64_t C;
36  int64_t M;
37  //storage_type* S;
38  std::atomic<T>* S;
39 
40  //T* S;
41 
42  explicit Array(int64_t c) :
43  C {c},
44  M {c-1},
45  //S {new storage_type[C]} {
46  //S {new T[static_cast<size_t>(C)]} {
47  S {new std::atomic<T>[static_cast<size_t>(C)]} {
48  //for(int64_t i=0; i<C; ++i) {
49  // ::new (std::addressof(S[i])) T();
50  //}
51  }
52 
53  ~Array() {
54  //for(int64_t i=0; i<C; ++i) {
55  // reinterpret_cast<T*>(std::addressof(S[i]))->~T();
56  //}
57  delete [] S;
58  }
59 
60  int64_t capacity() const noexcept {
61  return C;
62  }
63 
64  template <typename O>
65  void push(int64_t i, O&& o) noexcept {
66  //T* ptr = reinterpret_cast<T*>(std::addressof(S[i & M]));
67  //*ptr = std::forward<O>(o);
68  //S[i & M] = std::forward<O>(o);
69  S[i & M].store(std::forward<O>(o), std::memory_order_relaxed);
70  }
71 
72  T pop(int64_t i) noexcept {
73  //return *reinterpret_cast<T*>(std::addressof(S[i & M]));
74  //return S[i & M];
75  return S[i & M].load(std::memory_order_relaxed);
76  }
77 
78  Array* resize(int64_t b, int64_t t) {
79  Array* ptr = new Array {2*C};
80  for(int64_t i=t; i!=b; ++i) {
81  ptr->push(i, pop(i));
82  }
83  return ptr;
84  }
85 
86  };
87 
89  std::atomic<int64_t> _bottom;
90  std::atomic<Array*> _array;
91  std::vector<Array*> _garbage;
92  //char _padding[cacheline_size];
93 
94  public:
95 
101  explicit WorkStealingQueue(int64_t capacity = 1024);
102 
107 
111  bool empty() const noexcept;
112 
116  size_t size() const noexcept;
117 
121  int64_t capacity() const noexcept;
122 
134  template <typename O>
135  void push(O&& item);
136 
143  std::optional<T> pop();
144 
151  std::optional<T> steal();
152 };
153 
154 // Constructor
155 template <typename T>
157  assert(c && (!(c & (c-1))));
158  _top.store(0, std::memory_order_relaxed);
159  _bottom.store(0, std::memory_order_relaxed);
160  _array.store(new Array{c}, std::memory_order_relaxed);
161  _garbage.reserve(32);
162 }
163 
164 // Destructor
165 template <typename T>
167  for(auto a : _garbage) {
168  delete a;
169  }
170  delete _array.load();
171 }
172 
173 // Function: empty
174 template <typename T>
175 bool WorkStealingQueue<T>::empty() const noexcept {
176  int64_t b = _bottom.load(std::memory_order_relaxed);
177  int64_t t = _top.load(std::memory_order_relaxed);
178  return b <= t;
179 }
180 
181 // Function: size
182 template <typename T>
183 size_t WorkStealingQueue<T>::size() const noexcept {
184  int64_t b = _bottom.load(std::memory_order_relaxed);
185  int64_t t = _top.load(std::memory_order_relaxed);
186  return static_cast<size_t>(b >= t ? b - t : 0);
187 }
188 
189 // Function: push
190 template <typename T>
191 template <typename O>
193  int64_t b = _bottom.load(std::memory_order_relaxed);
194  int64_t t = _top.load(std::memory_order_acquire);
195  Array* a = _array.load(std::memory_order_relaxed);
196 
197  // queue is full
198  if(a->capacity() - 1 < (b - t)) {
199  Array* tmp = a->resize(b, t);
200  _garbage.push_back(a);
201  std::swap(a, tmp);
202  _array.store(a, std::memory_order_relaxed);
203  }
204 
205  a->push(b, std::forward<O>(o));
206  std::atomic_thread_fence(std::memory_order_release);
207  _bottom.store(b + 1, std::memory_order_relaxed);
208 }
209 
210 // Function: pop
211 template <typename T>
212 std::optional<T> WorkStealingQueue<T>::pop() {
213  int64_t b = _bottom.load(std::memory_order_relaxed) - 1;
214  Array* a = _array.load(std::memory_order_relaxed);
215  _bottom.store(b, std::memory_order_relaxed);
216  std::atomic_thread_fence(std::memory_order_seq_cst);
217  int64_t t = _top.load(std::memory_order_relaxed);
218 
219  std::optional<T> item;
220 
221  if(t <= b) {
222  item = a->pop(b);
223  if(t == b) {
224  // the last item just got stolen
225  if(!_top.compare_exchange_strong(t, t+1,
226  std::memory_order_seq_cst,
227  std::memory_order_relaxed)) {
228  item = std::nullopt;
229  }
230  _bottom.store(b + 1, std::memory_order_relaxed);
231  }
232  }
233  else {
234  _bottom.store(b + 1, std::memory_order_relaxed);
235  }
236 
237  return item;
238 }
239 
240 // Function: steal
241 template <typename T>
242 std::optional<T> WorkStealingQueue<T>::steal() {
243  int64_t t = _top.load(std::memory_order_acquire);
244  std::atomic_thread_fence(std::memory_order_seq_cst);
245  int64_t b = _bottom.load(std::memory_order_acquire);
246 
247  std::optional<T> item;
248 
249  if(t < b) {
250  Array* a = _array.load(std::memory_order_consume);
251  item = a->pop(t);
252  if(!_top.compare_exchange_strong(t, t+1,
253  std::memory_order_seq_cst,
254  std::memory_order_relaxed)) {
255  return std::nullopt;
256  }
257  }
258 
259  return item;
260 }
261 
262 // Function: capacity
263 template <typename T>
264 int64_t WorkStealingQueue<T>::capacity() const noexcept {
265  return _array.load(std::memory_order_relaxed)->capacity();
266 }
267 
268 } // end of namespace tf -----------------------------------------------------
bool empty() const noexcept
queries if the queue is empty at the time of this call
Definition: wsq.hpp:175
~WorkStealingQueue()
destructs the queue
Definition: wsq.hpp:166
void push(O &&item)
inserts an item to the queue
Definition: wsq.hpp:192
Definition: taskflow.hpp:5
size_t size() const noexcept
queries the number of items at the time of this call
Definition: wsq.hpp:183
int64_t capacity() const noexcept
queries the capacity of the queue
Definition: wsq.hpp:264
std::optional< T > pop()
pops out an item from the queue
Definition: wsq.hpp:212
Lock-free unbounded single-producer multiple-consumer queue.
Definition: wsq.hpp:27
std::optional< T > steal()
steals an item from the queue
Definition: wsq.hpp:242
WorkStealingQueue(int64_t capacity=1024)
constructs the queue with a given capacity
Definition: wsq.hpp:156