F´ Flight Software - C/C++ Documentation
A framework for building embedded system applications to NASA flight quality standards.
AtomicQueue.cpp
Go to the documentation of this file.
1 // ======================================================================
2 // \title AtomicQueue.cpp
3 // \author B. Duckett
4 // \brief Lock-free MPMC circular buffer with embedded buffer storage
5 //
6 // \copyright
7 // Copyright 2026, by the California Institute of Technology.
8 // ALL RIGHTS RESERVED. United States Government Sponsorship
9 // acknowledged.
10 //
11 // ======================================================================
12 
13 #include <Fw/Types/Assert.hpp>
15 #include <cstdio>
16 #include <cstring>
17 #include <new>
18 
19 namespace Types {
20 
21 U32 AtomicQueue::computeChecksum(const U8* buffer, FwSizeType size) {
22  FW_ASSERT(buffer != nullptr);
23  U32 sum = 0;
24  for (FwSizeType i = 0; i < size; ++i) {
25  sum += buffer[i];
26  sum = (sum << 1) | (sum >> 31); // Rotate left
27  }
28  return sum;
29 }
30 
32  : m_slots(nullptr),
33  m_bufferMemory(nullptr),
34  m_capacity(0),
35  m_bufferSize(0),
36  m_mask(0),
37  m_enqueuePos(0),
38  m_dequeuePos(0),
39  m_allocator(nullptr),
40  m_allocatorId(0),
41  m_notFullSem(nullptr) {}
42 
44  this->teardown();
45 }
46 
48  FwSizeType bufferSize,
49  Fw::MemAllocator& allocator,
50  FwEnumStoreType allocatorId) {
51  FW_ASSERT(numBuffers > 0, static_cast<FwAssertArgType>(numBuffers));
52  FW_ASSERT(bufferSize > 0, static_cast<FwAssertArgType>(bufferSize));
53 
54  this->m_capacity = numBuffers;
55  this->m_bufferSize = bufferSize;
56  this->m_allocator = &allocator;
57  this->m_allocatorId = allocatorId;
58 
59  // Optimization: use bitwise AND for power-of-2, otherwise modulo
60  bool isPowerOf2 = (numBuffers & (numBuffers - 1)) == 0;
61  this->m_mask = isPowerOf2 ? (numBuffers - 1) : 0;
62 
63  // Allocate slot array (with overflow check)
64  FW_ASSERT(numBuffers <= std::numeric_limits<FwSizeType>::max() / sizeof(Slot),
65  static_cast<FwAssertArgType>(numBuffers), static_cast<FwAssertArgType>(sizeof(Slot)));
66  FwSizeType slotsSize = numBuffers * sizeof(Slot);
67  void* slotMem = allocator.checkedAllocate(allocatorId, slotsSize, alignof(Slot));
68  FW_ASSERT(slotMem != nullptr, static_cast<FwAssertArgType>(numBuffers), static_cast<FwAssertArgType>(bufferSize));
69  this->m_slots = static_cast<Slot*>(slotMem);
70 
71  // Allocate contiguous buffer memory for all slots (with overflow check)
72  FW_ASSERT(numBuffers <= std::numeric_limits<FwSizeType>::max() / bufferSize,
73  static_cast<FwAssertArgType>(numBuffers), static_cast<FwAssertArgType>(bufferSize));
74  FwSizeType totalBufferSize = numBuffers * bufferSize;
75  void* bufferMem = allocator.checkedAllocate(allocatorId, totalBufferSize, 64);
76  FW_ASSERT(bufferMem != nullptr, static_cast<FwAssertArgType>(numBuffers), static_cast<FwAssertArgType>(bufferSize));
77  this->m_bufferMemory = static_cast<U8*>(bufferMem);
78 
79  // Initialize all slots with placement new and assign buffer pointers
80  for (FwSizeType i = 0; i < numBuffers; ++i) {
81  Slot* slot = new (&this->m_slots[i]) Slot();
82 
83  // Assign buffer from contiguous memory block
84  slot->buffer = this->m_bufferMemory + (i * bufferSize);
85  slot->size = 0;
86  slot->sequence.store(i, std::memory_order_relaxed);
87 
88  // Runtime verification that sequence atomics are lock-free
89  // This is critical for ISR safety and lock-free guarantee
90  FW_ASSERT(slot->sequence.is_lock_free(), static_cast<FwAssertArgType>(i),
91  static_cast<FwAssertArgType>(numBuffers));
92  }
93 
94  // Create semaphore for blocking enqueue support (all platforms)
95  // Allocate semaphore using provided allocator
96  FwSizeType semSize = sizeof(Os::CountingSemaphore);
97  void* semMem = allocator.checkedAllocate(allocatorId, semSize, alignof(Os::CountingSemaphore));
98  FW_ASSERT(semMem != nullptr, static_cast<FwAssertArgType>(numBuffers));
99 
100  // Use placement new to construct semaphore with initial count = numBuffers (all slots available)
101  this->m_notFullSem = new (semMem) Os::CountingSemaphore(static_cast<U32>(numBuffers));
102  FW_ASSERT(this->m_notFullSem != nullptr, static_cast<FwAssertArgType>(numBuffers));
103 
104  this->m_enqueuePos.store(0, std::memory_order_relaxed);
105  this->m_dequeuePos.store(0, std::memory_order_relaxed);
106 }
107 
109  // Destroy and deallocate semaphore
110  if (this->m_notFullSem != nullptr) {
111  FW_ASSERT(this->m_allocator != nullptr, 0);
112 
113  // Call destructor
114  this->m_notFullSem->~CountingSemaphore();
115 
116  // Deallocate memory
117  this->m_allocator->deallocate(this->m_allocatorId, this->m_notFullSem);
118  this->m_notFullSem = nullptr;
119  }
120 
121  // Destroy slots and deallocate memory
122  if (this->m_slots != nullptr) {
123  FW_ASSERT(this->m_capacity > 0, static_cast<FwAssertArgType>(this->m_capacity));
124  FW_ASSERT(this->m_allocator != nullptr, 0);
125 
126  // Call destructors on slots
127  for (FwSizeType i = 0; i < this->m_capacity; ++i) {
128  FW_ASSERT(i < this->m_capacity, static_cast<FwAssertArgType>(i),
129  static_cast<FwAssertArgType>(this->m_capacity));
130  this->m_slots[i].~Slot();
131  }
132 
133  // Deallocate buffer memory
134  if (this->m_bufferMemory != nullptr) {
135  this->m_allocator->deallocate(this->m_allocatorId, this->m_bufferMemory);
136  this->m_bufferMemory = nullptr;
137  }
138 
139  // Deallocate slot array
140  this->m_allocator->deallocate(this->m_allocatorId, this->m_slots);
141  this->m_slots = nullptr;
142  }
143 
144  this->m_enqueuePos.store(0, std::memory_order_relaxed);
145  this->m_dequeuePos.store(0, std::memory_order_relaxed);
146  this->m_capacity = 0;
147  this->m_bufferSize = 0;
148  this->m_mask = 0;
149  this->m_allocator = nullptr;
150 }
151 
152 bool AtomicQueue::enqueueInternal(const U8* buffer, FwSizeType size) {
153  FW_ASSERT(this->m_slots != nullptr, 0);
154  FW_ASSERT(buffer != nullptr, 0);
155  FW_ASSERT(size > 0, static_cast<FwAssertArgType>(size));
156  FW_ASSERT(size <= this->m_bufferSize, static_cast<FwAssertArgType>(size),
157  static_cast<FwAssertArgType>(this->m_bufferSize));
158 
159  for (FwSizeType retry = 0; retry < MAX_CAS_RETRIES; ++retry) {
160  FW_ASSERT(retry < MAX_CAS_RETRIES, static_cast<FwAssertArgType>(retry));
161 
162  // acquire-release on slot->sequence ensures coherence, enqueuePos & dequeuePos are relaxed since strong memory
163  // ordering is not required
164 
165  // Get current enqueue position
166  FwSizeType pos = this->m_enqueuePos.load(std::memory_order_relaxed);
167 
168  // Check against dequeue position to prevent lapping (detect full queue)
169  FwSizeType deqPos = this->m_dequeuePos.load(std::memory_order_relaxed);
170  FwSignedSizeType queueDiff = static_cast<FwSignedSizeType>(pos) - static_cast<FwSignedSizeType>(deqPos);
171  if (queueDiff >= static_cast<FwSignedSizeType>(this->m_capacity)) {
172  return false; // Queue is full
173  }
174 
175  Slot* slot = &this->m_slots[this->getIndex(pos)];
176  FwSizeType seq = slot->sequence.load(std::memory_order_acquire);
177 
178  // Check if slot is ready for write (seq == pos means available)
179  FwSignedSizeType diff = static_cast<FwSignedSizeType>(seq) - static_cast<FwSignedSizeType>(pos);
180 
181  if (diff == 0) {
182  // Slot available, try to claim it
183  if (this->m_enqueuePos.compare_exchange_weak(pos, pos + 1, std::memory_order_release,
184  std::memory_order_relaxed)) {
185  // Claimed the slot, copy message data
186  FW_ASSERT(slot->buffer != nullptr, static_cast<FwAssertArgType>(pos));
187  (void)std::memcpy(slot->buffer, buffer, size);
188  slot->size = size;
189 
190  // Mark slot as ready for read
191  slot->sequence.store(pos + 1, std::memory_order_release);
192  return true;
193  }
194  } else if (diff < 0) {
195  // Queue is full (wrapped around)
196  return false;
197  }
198  // else: another producer claimed this slot, retry
199  }
200 
201  return false;
202 }
203 
204 bool AtomicQueue::enqueue(const U8* buffer, FwSizeType size) {
205  bool success = this->enqueueInternal(buffer, size);
206 
207  // Decrement semaphore to track available slots (if semaphore exists)
208  // This ensures blocking sends see correct availability even when
209  // queue is filled via non-blocking sends.
210  //
211  // NOTE: tryWait() may fail due to race conditions when multiple threads
212  // concurrently enqueue. This is acceptable - the semaphore is a best-effort
213  // hint for blocking operations. The lock-free atomics in enqueueInternal()
214  // are the authoritative source of queue state.
215  if (success && this->m_notFullSem != nullptr) {
216  (void)this->m_notFullSem->tryWait();
217  }
218 
219  return success;
220 }
221 
222 bool AtomicQueue::enqueueBlocking(const U8* buffer, FwSizeType size, bool blockIfFull) {
223  FW_ASSERT(this->m_slots != nullptr, 0);
224  FW_ASSERT(buffer != nullptr, 0);
225 
226  // If not blocking or no semaphore, just use non-blocking enqueue
227  if (!blockIfFull || this->m_notFullSem == nullptr) {
228  return this->enqueue(buffer, size);
229  }
230 
231  // Wait-first pattern: reserve slot via semaphore, then enqueue
232  // This prevents semaphore count drift under contention
233  for (FwSizeType attempt = 0; attempt < MAX_CAS_RETRIES; ++attempt) {
234  FW_ASSERT(attempt < MAX_CAS_RETRIES, static_cast<FwAssertArgType>(attempt));
235 
236  // Reserve a slot by waiting on semaphore (blocks until space available)
237  Os::CountingSemaphoreInterface::Status status = this->m_notFullSem->wait();
238  FW_ASSERT(status == Os::CountingSemaphoreInterface::Status::OP_OK, static_cast<FwAssertArgType>(status));
239 
240  // Try to enqueue (should succeed since we reserved a slot)
241  // Use internal method to avoid double-decrementing semaphore
242  if (this->enqueueInternal(buffer, size)) {
243  return true; // Success
244  }
245 
246  // Extremely rare: slot was stolen between wait() and enqueue()
247  // Return the semaphore permit and retry
248  status = this->m_notFullSem->post();
249  FW_ASSERT(status == Os::CountingSemaphoreInterface::Status::OP_OK, static_cast<FwAssertArgType>(status));
250  }
251 
252  // Exceeded retry limit (should never happen in practice)
253  return false;
254 }
255 
256 bool AtomicQueue::dequeue(U8* buffer, FwSizeType capacity, FwSizeType& actualSize) {
257  FW_ASSERT(this->m_slots != nullptr, 0);
258  FW_ASSERT(buffer != nullptr, 0);
259  FW_ASSERT(capacity > 0, static_cast<FwAssertArgType>(capacity));
260 
261  for (FwSizeType retry = 0; retry < MAX_CAS_RETRIES; ++retry) {
262  FW_ASSERT(retry < MAX_CAS_RETRIES, static_cast<FwAssertArgType>(retry));
263 
264  // acquire-release on slot->sequence ensures coherence, enqueuePos & dequeuePos are relaxed since strong memory
265  // ordering is not required
266 
267  // Get current dequeue & enqueue positions
268  FwSizeType pos = this->m_dequeuePos.load(std::memory_order_relaxed);
269  Slot* slot = &this->m_slots[this->getIndex(pos)];
270  FwSizeType seq = slot->sequence.load(std::memory_order_acquire);
271 
272  // Check if slot is ready for read (seq == pos + 1 means data available)
273  FwSignedSizeType diff = static_cast<FwSignedSizeType>(seq) - static_cast<FwSignedSizeType>(pos + 1);
274 
275  if (diff == 0) {
276  // Slot has data, try to claim it
277  if (this->m_dequeuePos.compare_exchange_weak(pos, pos + 1, std::memory_order_release,
278  std::memory_order_relaxed)) {
279  // Claimed the slot, check size and copy data
280  // Note: sequence.load(acquire) at 10 lines above already synchronizes with
281  // enqueue's sequence.store(release), ensuring size & buffer coherency
282  FW_ASSERT(slot->buffer != nullptr, static_cast<FwAssertArgType>(pos));
283  FW_ASSERT(slot->size > 0, static_cast<FwAssertArgType>(slot->size));
284  FW_ASSERT(slot->size <= capacity, static_cast<FwAssertArgType>(slot->size),
285  static_cast<FwAssertArgType>(capacity));
286 
287  actualSize = slot->size;
288  (void)std::memcpy(buffer, slot->buffer, actualSize);
289 
290  // Mark slot as available for next cycle (pos + capacity)
291  slot->sequence.store(pos + this->m_capacity, std::memory_order_release);
292 
293  // Post semaphore to wake up blocked enqueue threads (if semaphore exists)
294  // ISR-SAFETY NOTE: Calling from ISR depends on platform semaphore implementation.
295  // See header comments for details.
296  if (this->m_notFullSem != nullptr) {
297  Os::CountingSemaphoreInterface::Status status = this->m_notFullSem->post();
299  static_cast<FwAssertArgType>(status));
300  }
301 
302  return true;
303  }
304  } else if (diff < 0) {
305  // Queue is empty (no data written yet)
306  return false;
307  }
308  // else: another consumer claimed this slot, retry
309  }
310 
311  return false;
312 }
313 
314 bool AtomicQueue::isFull() const {
315  // Queue is full if enqueue is exactly capacity ahead of dequeue
316  // Return false for uninitialized queue
317  if (this->m_capacity == 0) {
318  return false;
319  }
320  return this->getSize() >= this->m_capacity;
321 }
322 
323 bool AtomicQueue::isEmpty() const {
324  return this->getSize() == 0;
325 }
326 
328  // Safe to call on uninitialized queue
329  if (this->m_capacity == 0) {
330  return 0;
331  }
332  // A nonzero capacity implies the queue was created with backing storage
333  FW_ASSERT(this->m_slots != nullptr);
334 
335  FwSizeType enq = this->m_enqueuePos.load(std::memory_order_relaxed);
336  FwSizeType deq = this->m_dequeuePos.load(std::memory_order_relaxed);
337  FwSignedSizeType diff = static_cast<FwSignedSizeType>(enq) - static_cast<FwSignedSizeType>(deq);
338 
339  // Two independent relaxed loads provide no cross-variable consistency guarantee.
340  // Restrict to [0, capacity] rather than asserting — diff can be transiently negative
341  // or > capacity on concurrent access even though no real program state has that.
342  if (diff < 0) {
343  return 0;
344  }
345  if (static_cast<FwSizeType>(diff) > this->m_capacity) {
346  return this->m_capacity;
347  }
348  return static_cast<FwSizeType>(diff);
349 }
350 
352  // Safe to call on uninitialized queue - returns 0 if not created
353  return this->m_capacity;
354 }
355 
357  // Safe to call on uninitialized queue - returns 0 if not created
358  return this->m_bufferSize;
359 }
360 
361 } // namespace Types
~CountingSemaphore() final
Destructor.
Operation succeeded.
Definition: Os.hpp:27
PlatformSizeType FwSizeType
I32 FwEnumStoreType
~AtomicQueue()
AtomicQueue destructor.
Definition: AtomicQueue.cpp:43
AtomicQueue()
AtomicQueue constructor.
Definition: AtomicQueue.cpp:31
PlatformSignedSizeType FwSignedSizeType
FwSizeType getSize() const
Get the current number of elements in the queue.
bool enqueueBlocking(const U8 *buffer, FwSizeType size, bool blockIfFull)
Enqueue with optional blocking (multi-producer safe, O(1))
Status tryWait() override
non-blocking attempt to decrement the semaphore
bool isEmpty() const
Check if the queue is empty.
void * checkedAllocate(const FwEnumStoreType identifier, FwSizeType &size, bool &recoverable, FwSizeType alignment=alignof(std::max_align_t))
bool isFull() const
Check if the queue is full.
uint8_t U8
8-bit unsigned integer
Definition: BasicTypes.h:54
Memory Allocation base class.
void teardown()
Teardown the queue and free allocated memory.
FwSizeType getBufferSize() const
Get the buffer size for each message.
bool enqueue(const U8 *buffer, FwSizeType size)
Enqueue a message (multi-producer safe, non-blocking, O(1))
FwSizeType getCapacity() const
Get the maximum capacity of the queue.
Status wait() override
wait (decrement), blocking if count is zero
Status post() override
post (increment) the semaphore, potentially waking a waiting thread
void create(FwSizeType numBuffers, FwSizeType bufferSize, Fw::MemAllocator &allocator, FwEnumStoreType allocatorId)
Create the queue with embedded buffer storage.
Definition: AtomicQueue.cpp:47
virtual void deallocate(const FwEnumStoreType identifier, void *ptr)=0
bool dequeue(U8 *buffer, FwSizeType capacity, FwSizeType &actualSize)
Dequeue a message (multi-consumer safe, non-blocking, O(1))
#define FW_ASSERT(...)
Definition: Assert.hpp:14
PlatformAssertArgType FwAssertArgType
The type of arguments to assert functions.