arduino-audio-tools
Loading...
Searching...
No Matches
QueueLockFree.h
Go to the documentation of this file.
1
2#pragma once
3#include <stdint.h>
4
5#include <atomic>
6#include <cstddef>
7#include <utility>
8
10
11namespace audio_tools {
12
21template <typename T>
23 public:
25 setAllocator(allocator);
27 }
28
30 size_t head = head_pos.load(std::memory_order_acquire);
31 size_t tail = tail_pos.load(std::memory_order_acquire);
32 for (size_t i = head; i != tail; ++i)
33 p_node[i & capacity_mask].ptr()->~T();
34 }
35
36 void setAllocator(Allocator& allocator) { vector.setAllocator(allocator); }
37
38 bool resize(size_t capacity) {
39 if (capacity == 0) capacity = 1;
40
41 // Destroy any live elements in the current queue before reinitialising.
42 if (p_node) {
43 size_t head = head_pos.load(std::memory_order_relaxed);
44 size_t tail = tail_pos.load(std::memory_order_relaxed);
45 for (size_t i = head; i != tail; ++i)
46 p_node[i & capacity_mask].ptr()->~T();
47 }
48
49 // Round capacity up to the next power of two so that the bitmask index
50 // wrapping always stays within the allocated array.
51 size_t new_capacity_mask = capacity - 1;
52 for (size_t i = 1; i <= sizeof(void*) * 4; i <<= 1)
53 new_capacity_mask |= new_capacity_mask >> i;
54 size_t new_capacity_value = new_capacity_mask + 1;
55
56 vector.resize(new_capacity_value);
57 p_node = vector.data();
58 if (p_node == nullptr) {
59 // Allocation failed: leave the queue empty rather than dereferencing
60 // a null p_node below.
61 capacity_mask = 0;
63 tail_pos.store(0, std::memory_order_relaxed);
64 head_pos.store(0, std::memory_order_relaxed);
65 return false;
66 }
67 capacity_mask = new_capacity_mask;
68 capacity_value = new_capacity_value;
69
70 for (size_t i = 0; i < capacity_value; ++i) {
71 p_node[i].tail.store(i, std::memory_order_relaxed);
72 p_node[i].head.store(size_t(-1), std::memory_order_relaxed);
73 }
74
75 tail_pos.store(0, std::memory_order_relaxed);
76 head_pos.store(0, std::memory_order_relaxed);
77 return true;
78 }
79
80 size_t capacity() const { return capacity_value; }
81
82 size_t availableForWrite() const {
83 size_t head = head_pos.load(std::memory_order_seq_cst);
84 size_t tail = tail_pos.load(std::memory_order_seq_cst);
85 return capacity_value - (tail - head);
86 }
87
88 bool empty() const { return size() == 0; }
89
90 size_t size() const {
91 size_t head = head_pos.load(std::memory_order_seq_cst);
92 size_t tail = tail_pos.load(std::memory_order_seq_cst);
93 return tail - head;
94 }
95
96 bool enqueue(const T& data) { return emplace(data); }
97 bool enqueue(T&& data) { return emplace(std::move(data)); }
98
99 bool dequeue(T& result) {
100 if (capacity_value == 0) return false;
101 Node* node;
102 size_t head = head_pos.load(std::memory_order_relaxed);
103 for (;;) {
104 node = &p_node[head & capacity_mask];
105 // acquire pairs with enqueue's release store on node->head,
106 // ensuring node->storage is visible before we read it.
107 if (node->head.load(std::memory_order_acquire) != head) return false;
108 if (head_pos.compare_exchange_weak(head, head + 1,
109 std::memory_order_relaxed))
110 break;
111 }
112 result = std::move(*node->ptr());
113 node->ptr()->~T();
114 node->tail.store(head + capacity_value, std::memory_order_release);
115 return true;
116 }
117
118 void clear() {
119 size_t head = head_pos.load(std::memory_order_acquire);
120 size_t tail = tail_pos.load(std::memory_order_acquire);
121 for (size_t i = head; i != tail; ++i) {
122 Node* node = &p_node[i & capacity_mask];
123 node->ptr()->~T();
124 node->tail.store(i + capacity_value, std::memory_order_release);
125 }
126 head_pos.store(tail, std::memory_order_release);
127 }
128
129 protected:
130 struct Node {
131 alignas(T) unsigned char storage[sizeof(T)];
132 std::atomic<size_t> tail;
133 std::atomic<size_t> head;
134 T* ptr() { return reinterpret_cast<T*>(storage); }
135 const T* ptr() const { return reinterpret_cast<const T*>(storage); }
136 };
137
138 // Single enqueue implementation for both lvalue and rvalue paths.
139 template <typename U>
140 bool emplace(U&& val) {
141 if (capacity_value == 0) return false;
142 Node* node;
143 size_t tail = tail_pos.load(std::memory_order_relaxed);
144 for (;;) {
145 node = &p_node[tail & capacity_mask];
146 // acquire pairs with dequeue's release store on node->tail,
147 // ensuring ~T() in the consumer is complete before we reuse the slot.
148 if (node->tail.load(std::memory_order_acquire) != tail) return false;
149 if (tail_pos.compare_exchange_weak(tail, tail + 1,
150 std::memory_order_relaxed))
151 break;
152 }
153 new (node->ptr()) T(std::forward<U>(val));
154 node->head.store(tail, std::memory_order_release);
155 return true;
156 }
157
158 Node* p_node = nullptr;
159 size_t capacity_mask = 0;
160 size_t capacity_value = 0;
161 std::atomic<size_t> tail_pos{0};
162 std::atomic<size_t> head_pos{0};
164};
165
166} // namespace audio_tools
Memory allocateator which uses malloc.
Definition Allocator.h:25
Lock-free MPMC queue.
Definition QueueLockFree.h:22
size_t availableForWrite() const
Definition QueueLockFree.h:82
bool enqueue(T &&data)
Definition QueueLockFree.h:97
bool dequeue(T &result)
Definition QueueLockFree.h:99
size_t size() const
Definition QueueLockFree.h:90
size_t capacity_value
Definition QueueLockFree.h:160
bool resize(size_t capacity)
Definition QueueLockFree.h:38
bool empty() const
Definition QueueLockFree.h:88
size_t capacity() const
Definition QueueLockFree.h:80
bool enqueue(const T &data)
Definition QueueLockFree.h:96
std::atomic< size_t > tail_pos
Definition QueueLockFree.h:161
bool emplace(U &&val)
Definition QueueLockFree.h:140
void clear()
Definition QueueLockFree.h:118
void setAllocator(Allocator &allocator)
Definition QueueLockFree.h:36
std::atomic< size_t > head_pos
Definition QueueLockFree.h:162
~QueueLockFree()
Definition QueueLockFree.h:29
size_t capacity_mask
Definition QueueLockFree.h:159
QueueLockFree(size_t capacity, Allocator &allocator=DefaultAllocator)
Definition QueueLockFree.h:24
Vector< Node > vector
Definition QueueLockFree.h:163
Node * p_node
Definition QueueLockFree.h:158
Vector implementation which provides the most important methods as defined by std::vector....
Definition Vector.h:21
Generic Implementation of sound input and output for desktop environments using portaudio.
Definition LMSEchoCancellationStream.h:6
static TAllocatorExt DefaultAllocator
Definition Allocator.h:208
Definition QueueLockFree.h:130
unsigned char storage[sizeof(T)]
Definition QueueLockFree.h:131
const T * ptr() const
Definition QueueLockFree.h:135
std::atomic< size_t > head
Definition QueueLockFree.h:133
T * ptr()
Definition QueueLockFree.h:134
std::atomic< size_t > tail
Definition QueueLockFree.h:132