Aleph-w 3.0
A C++ Library for Data Structures and Algorithms
Loading...
Searching...
No Matches
tpl_mpsc_queue.H
Go to the documentation of this file.
1
2/*
3 Aleph_w
4
5 Data structures & Algorithms
6 https://github.com/lrleon/Aleph-w
7
8 This file is part of Aleph-w library
9
10 Copyright (c) 2002-2026 Leandro Rabindranath Leon
11
12 Permission is hereby granted, free of charge, to any person obtaining a copy
13 of this software and associated documentation files (the "Software"), to deal
14 in the Software without restriction, including without limitation the rights
15 to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
16 copies of the Software, and to permit persons to whom the Software is
17 furnished to do so, subject to the following conditions:
18
19 The above copyright notice and this permission notice shall be included in all
20 copies or substantial portions of the Software.
21
22 THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
23 IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
24 FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
25 AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
26 LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
27 OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
28 SOFTWARE.
29*/
30
106#ifndef TPL_MPSC_QUEUE_H
107#define TPL_MPSC_QUEUE_H
108
109#include <atomic>
110#include <memory>
111#include <optional>
112#include <utility>
113
114namespace Aleph
115{
116
126template <typename T>
128{
129 struct Node
130 {
131 std::atomic<Node *> next{nullptr};
132 std::optional<T> value;
133
134 Node() = default;
135
136 template <typename... Args>
137 explicit Node(std::in_place_t, Args &&... args)
138 : value(std::in_place, std::forward<Args>(args)...)
139 {}
140 };
141
142 alignas(64) std::atomic<Node *> head_;
143 alignas(64) Node * tail_;
145
148 void push_node(Node * n) noexcept
149 {
150 n->next.store(nullptr, std::memory_order_relaxed);
151 Node * prev = head_.exchange(n, std::memory_order_acq_rel);
152 prev->next.store(n, std::memory_order_release);
153 }
154
166 {
167 Node * tail = tail_;
168 Node * next = tail->next.load(std::memory_order_acquire);
169
170 if (tail == &stub_)
171 {
172 if (next == nullptr)
173 return nullptr; // genuinely empty
174 tail_ = next;
175 tail = next;
176 next = next->next.load(std::memory_order_acquire);
177 }
178
179 if (next != nullptr)
180 {
181 tail_ = next;
182 return tail;
183 }
184
185 Node * head = head_.load(std::memory_order_acquire);
186 if (tail != head)
187 return nullptr; // a producer is mid-push; try again later
188
189 push_node(&stub_); // re-anchor the queue so progress remains possible
190 next = tail->next.load(std::memory_order_acquire);
191 if (next != nullptr)
192 {
193 tail_ = next;
194 return tail;
195 }
196
197 return nullptr;
198 }
199
200public:
208
215 {
216 // stub_ is a data member, not a heap node: it must never reach
217 // `delete`. A racing push_node(&stub_) re-anchor (see pop_node()) can
218 // leave stub_ reachable from tail_ without *being* tail_ (a real,
219 // not-yet-popped node can end up between them), so every node in the
220 // remaining chain -- not just the first -- must be checked before
221 // deleting it.
222 Node * n = tail_;
223 while (n != nullptr)
224 {
225 Node * next = n->next.load(std::memory_order_relaxed);
226 if (n != &stub_)
227 delete n;
228 n = next;
229 }
230 }
231
234 MpscQueue(const MpscQueue &) = delete;
236 MpscQueue & operator=(const MpscQueue &) = delete;
239 MpscQueue(MpscQueue &&) = delete;
242
252 void push(const T & value) { emplace(value); }
253
264 void push(T && value) { emplace(std::move(value)); }
265
277 template <typename... Args>
278 void emplace(Args &&... args)
279 {
280 Node * n = new Node(std::in_place, std::forward<Args>(args)...);
281 push_node(n);
282 }
283
298 bool try_pop(T & out)
299 {
300 Node * n = pop_node();
301 if (n == nullptr)
302 return false;
303 // unique_ptr guarantees the node is freed even if the move below
304 // throws (T's move assignment is not required to be noexcept).
305 std::unique_ptr<Node> owned(n);
306 out = std::move(*owned->value);
307 return true;
308 }
309
322 [[nodiscard]] std::optional<T> try_pop()
323 {
324 Node * n = pop_node();
325 if (n == nullptr)
326 return std::nullopt;
327 std::unique_ptr<Node> owned(n);
328 return std::optional<T>(std::move(*owned->value));
329 }
330
342 {
343 return tail_ == &stub_ and
344 tail_->next.load(std::memory_order_acquire) == nullptr;
345 }
346};
347
353template <typename T>
355
356} // namespace Aleph
357
358#endif // TPL_MPSC_QUEUE_H
size_t size_t int32_t value
Definition ca-c-api.h:116
size_t size_t int32_t * out
Definition ca-c-api.h:120
Unbounded lock-free multi-producer/single-consumer queue.
void push_node(Node *n) noexcept
Publish n as the new last node.
std::atomic< Node * > head_
MpscQueue() noexcept
Construct an empty queue.
~MpscQueue()
Destroy the queue, releasing any still-queued nodes.
void emplace(Args &&... args)
Construct a new element in place at the back of the queue.
void push(const T &value)
Push a copy of value onto the queue.
void push(T &&value)
Push value onto the queue, moving it in.
MpscQueue & operator=(const MpscQueue &)=delete
Deleted copy assignment operator.
MpscQueue(const MpscQueue &)=delete
Deleted copy constructor: the queue owns heap nodes with internal atomics that cannot be safely dupli...
std::optional< T > try_pop()
Attempt to pop the front element.
MpscQueue & operator=(MpscQueue &&)=delete
Deleted move assignment operator.
Node * pop_node() noexcept
Consumer-only: detach and return the front node, or nullptr.
MpscQueue(MpscQueue &&)=delete
Deleted move constructor: producers may hold a reference to a fixed queue address; see the class-leve...
bool is_empty() const noexcept
Advisory check for whether the queue currently has no elements.
bool try_pop(T &out)
Attempt to pop the front element into out.
size_t blossom_maximum_cardinality_matching(const GT &g, DynDlist< typename GT::Arc * > &matching, SA sa=SA())
Alias of compute_maximum_cardinality_general_matching().
Definition Blossom.H:466
Main namespace for Aleph-w library functions.
Definition ah-arena.H:89
and
Check uniqueness with explicit hash + equality functors.
std::decay_t< typename HeadC::Item_Type > T
Definition ah-zip.H:105
void next()
Advance all underlying iterators (bounds-checked).
Definition ah-zip.H:171
STL namespace.
std::atomic< Node * > next
Node(std::in_place_t, Args &&... args)
std::optional< T > value