46 bool PriorityMemQueue::s_requirePrioritySizing =
false;
47 std::atomic<bool>* PriorityMemQueue::s_configsUsed =
nullptr;
48 bool PriorityMemQueue::s_configured =
false;
55 return 1U << priority;
97 if (static_cast<FwQueuePriorityType>(p) > this->
m_maxPriority) {
108 if (atomicQueuesMem ==
nullptr) {
119 FwSizeType hwmSize =
sizeof(std::atomic<U32>) * this->m_numActivePriorities;
120 void* hwmMem = allocator.
checkedAllocate(allocatorId, hwmSize,
alignof(std::atomic<U32>));
121 if (hwmMem ==
nullptr) {
164 FW_ASSERT(priority < Os::Generic::Queue::MAX_PRIORITIES, priority, this->
m_id);
174 FW_ASSERT(priority < Os::Generic::Queue::MAX_PRIORITIES, this->
m_id);
194 #if defined(__GNUC__) || defined(__clang__) 195 I32 msb = 31 - __builtin_clz(value);
200 U32 mask = 0x80000000;
217 if (priorities == 0) {
221 if (priorities == 0) {
257 for (
FwSizeType i = 0; i < numQueueConfigs; ++i) {
263 static_cast<FwAssertArgType>(i), currentConfig->
instanceId,
267 for (
FwSizeType j = i + 1; j < numQueueConfigs; ++j) {
270 static_cast<FwAssertArgType>(i), static_cast<FwAssertArgType>(j));
275 FW_ASSERT(priorityConfigs !=
nullptr, static_cast<FwAssertArgType>(i), currentConfig->
instanceId,
281 FW_ASSERT(pConfig->
maxMsgSize > 0, static_cast<FwAssertArgType>(i), static_cast<FwAssertArgType>(p),
283 FW_ASSERT(pConfig->
numMsgs > 0, static_cast<FwAssertArgType>(i), static_cast<FwAssertArgType>(p),
286 static_cast<FwAssertArgType>(i), static_cast<FwAssertArgType>(p), pConfig->
priority);
303 FW_ASSERT((queueConfigs !=
nullptr) || (numQueueConfigs == 0), 0);
311 FW_ASSERT(!(required && numQueueConfigs == 0), required, static_cast<FwAssertArgType>(numQueueConfigs));
314 if (queueConfigs !=
nullptr) {
321 if (numQueueConfigs > 0) {
322 FwSizeType expSize = numQueueConfigs *
sizeof(std::atomic<bool>);
323 s_configsUsed =
static_cast<std::atomic<bool>*
>(
324 allocator.
checkedAllocate(allocatorId, expSize,
alignof(std::atomic<bool>)));
327 for (
FwSizeType i = 0; i < numQueueConfigs; ++i) {
328 new (&s_configsUsed[i]) std::atomic<bool>(
false);
334 "QueueConfig array must be naturally aligned with QueuePriorityConfig");
335 if (numQueueConfigs > 0) {
337 for (
FwSizeType i = 0; i < numQueueConfigs; ++i) {
344 for (
FwSizeType i = 0; i < numQueueConfigs; ++i) {
345 s_configs[i] = queueConfigs[i];
347 (void)memcpy(priorityBase, queueConfigs[i].priorityConfigs, priorityConfigsSize);
352 s_numConfigs = numQueueConfigs;
353 s_requirePrioritySizing = required;
354 s_allocatorId = allocatorId;
359 FW_ASSERT((s_configs !=
nullptr) || (s_numConfigs == 0), static_cast<FwAssertArgType>(s_numConfigs));
361 if (s_configsUsed !=
nullptr || s_configs !=
nullptr) {
368 if (s_configsUsed !=
nullptr) {
369 allocator.
deallocate(allocatorId, s_configsUsed);
370 s_configsUsed =
nullptr;
374 if (s_configs !=
nullptr) {
382 s_requirePrioritySizing =
false;
383 s_configured =
false;
407 if (queueConfig !=
nullptr) {
422 return Os::QueueInterface::Status::ALLOCATION_FAILED;
431 if (semMem ==
nullptr) {
433 return Os::QueueInterface::Status::ALLOCATION_FAILED;
438 if (queueConfig !=
nullptr) {
440 return createConfiguredQueues(queueConfig, allocator,
id);
443 return createDefaultQueue(depth, messageSize, allocator,
id);
449 if (s_configs !=
nullptr && s_configsUsed !=
nullptr) {
450 for (
FwSizeType i = 0; i < s_numConfigs; ++i) {
451 if (s_configs[i].instanceId ==
id) {
453 bool expected =
false;
454 if (s_configsUsed[i].compare_exchange_strong(expected,
true, std::memory_order_acq_rel)) {
455 return &s_configs[i];
458 FW_ASSERT(
false,
id, s_configs[i].instanceId);
471 FW_ASSERT(queueConfig->priorityConfigs !=
nullptr, this->m_handle.m_id,
472 static_cast<FwAssertArgType>(queueConfig->numPriorities));
474 for (
FwSizeType i = 0; i < queueConfig->numPriorities; ++i) {
475 const QueuePriorityConfig& priorityConfig = queueConfig->priorityConfigs[i];
480 priorityConfig.numMsgs, allocator, allocatorId);
487 this->setPriorityEnabled(priority,
true);
520 atomicQueue->
create(numMsgs, maxMsgSize, allocator, allocatorId);
525 return Os::QueueInterface::Status::ALLOCATION_FAILED;
529 this->setPriorityEnabled(priority,
true);
568 if (s_configs !=
nullptr && s_configsUsed !=
nullptr) {
569 for (
FwSizeType i = 0; i < s_numConfigs; ++i) {
570 if (s_configs[i].instanceId == this->
m_handle.
m_id && s_configsUsed[i].load()) {
571 s_configsUsed[i].store(
false);
587 bool requirePrioritySizing) {
593 if (requirePrioritySizing) {
598 FW_ASSERT(index >= 0, queueId, priority);
618 FW_ASSERT(highWaterMarks !=
nullptr, queueId, static_cast<FwAssertArgType>(index));
619 U32 prevMax = highWaterMarks[index].load(std::memory_order_acquire);
621 constexpr U32 MAX_CAS_RETRIES = 100;
622 for (U32 casRetries = 0; casRetries < MAX_CAS_RETRIES; ++casRetries) {
623 if (currentDepth <= prevMax) {
626 if (highWaterMarks[index].compare_exchange_weak(prevMax, currentDepth, std::memory_order_release,
627 std::memory_order_acquire)) {
643 return QueueInterface::Status::INVALID_PRIORITY;
648 return QueueInterface::Status::UNINITIALIZED;
657 return QueueInterface::Status::SIZE_MISMATCH;
662 if (blockType == QueueInterface::BlockingType::BLOCKING) {
665 success = atomicQueue->
enqueue(buffer, size);
669 return QueueInterface::Status::FULL;
675 U32 currentDepth =
static_cast<U32
>(atomicQueue->
getSize());
677 this->m_handle.m_id);
697 return QueueInterface::Status::UNINITIALIZED;
702 blockType == QueueInterface::BlockingType::BLOCKING || blockType == QueueInterface::BlockingType::NONBLOCKING,
731 const bool dequeued = aq->
dequeue(destination, capacity, actualSize);
733 priority = testPriority;
739 if (blockType == QueueInterface::BlockingType::BLOCKING) {
743 static_cast<FwAssertArgType>(semStatus));
745 return QueueInterface::Status::EMPTY;
751 return QueueInterface::Status::UNKNOWN_ERROR;
759 static_cast<FwAssertArgType>(this->m_handle.m_numActivePriorities));
763 total += atomicQueue->
getSize();
777 static_cast<FwAssertArgType>(this->m_handle.m_numActivePriorities));
~CountingSemaphore() final
Destructor.
FwEnumStoreType m_allocatorId
I8 m_priorityMap[Queue::MAX_PRIORITIES]
void enablePriority(FwQueuePriorityType priority)
Enable a specific priority (in the handle )
std::atomic< U32 > m_priorityMask
PlatformSizeType FwSizeType
FwQueuePriorityType priority
void teardown() override
teardown the queue
Status
status returned from the queue send function
static constexpr FwSizeType MAX_PRIORITIES
constexpr U32 LOOP_GUARD_LIMIT
Configuration for a priority level.
QueueHandle parent class.
int8_t I8
8-bit signed integer
static void configure(QueueConfig *queueConfigs, FwSizeType numQueueConfigs, bool required, FwEnumStoreType allocatorId)
Register per-priority sizes for component instances which do queueing with multiple priorities...
~AtomicQueue()
AtomicQueue destructor.
static constexpr FwSizeType DEFAULT_PRIORITY
FwSizeType getSize() const
Get the current number of elements in the queue.
static MemAllocatorRegistry & getInstance()
get the singleton registry
QueueHandle * getHandle() override
return the underlying queue handle (implementation specific)
bool enqueueBlocking(const U8 *buffer, FwSizeType size, bool blockIfFull)
Enqueue with optional blocking (multi-producer safe, O(1))
virtual ~PriorityMemQueue()
default queue destructor
Types::AtomicQueue * m_atomicQueues
PriorityMemQueueHandle m_handle
static void updateHighWaterMark(std::atomic< U32 > *highWaterMarks, FwSizeType index, U32 currentDepth, FwEnumStoreType queueId)
Update per-priority high water mark atomically.
A lock-free MPMC FIFO circular buffer with fixed-size buffer storage.
Status receive(U8 *destination, FwSizeType capacity, BlockingType blockType, FwSizeType &actualSize, FwQueuePriorityType &priority) override
receive a message from the queue
REQUIRED: required for Os::Queue memory allocation when using queues that allocate memory...
MemAllocator & getAnAllocator(const MemoryAllocation::MemoryAllocatorType type)
static Types::AtomicQueue * resolvePriorityQueue(PriorityMemQueueHandle &handle, FwQueuePriorityType &priority, FwEnumStoreType queueId, bool requirePrioritySizing)
Resolve priority to a valid AtomicQueue, fallback to DEFAULT if needed.
static void resetConfig()
Reset static configuration (test environments only)
FwSizeType getMessagesAvailable() const override
get number of messages available
void init()
Initialize the handle.
void * checkedAllocate(const FwEnumStoreType identifier, FwSizeType &size, bool &recoverable, FwSizeType alignment=alignof(std::max_align_t))
void deallocateArrays(Fw::MemAllocator &allocator, FwEnumStoreType allocatorId)
Deallocate arrays for priority data.
Status create(FwEnumStoreType id, const Fw::ConstStringBase &name, FwSizeType depth, FwSizeType messageSize) override
create queue storage
bool allocateArrays(Fw::MemAllocator &allocator, FwEnumStoreType allocatorId)
Allocate arrays for priority data (sparse allocation) Uses m_priorityMap to determine which prioritie...
bool isCreated() const
Check if queue has been successfully created.
Os::CountingSemaphore * m_notEmptySem
uint8_t U8
8-bit unsigned integer
Configuration for a queue with multiple priorities.
PlatformQueuePriorityType FwQueuePriorityType
The type of queue priorities used.
critical data stored for priority queue
static constexpr U32 priorityBitMask(FwQueuePriorityType priority)
Get the bit mask for a priority.
Memory Allocation base class.
I8 getPriorityIndex(FwQueuePriorityType priority) const
Get array index for a priority (inline for performance)
void teardown()
Teardown the queue and free allocated memory.
FwQueuePriorityType m_maxPriority
A read-only abstract superclass for StringBase.
void disablePriority(FwQueuePriorityType priority)
Disable a specific priority.
static void validateQueueConfigs(PriorityMemQueue::QueueConfig *queueConfigs, FwSizeType numQueueConfigs)
Validate queue configuration structures.
void teardownInternal()
teardown the queue
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))
Status wait() override
wait (decrement), blocking if count is zero
Defines a base class for a memory allocator for classes.
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.
static I32 findMSB(U32 value)
Find most significant bit set (IPC-style priority finding)
virtual void deallocate(const FwEnumStoreType identifier, void *ptr)=0
std::atomic< U32 > * m_highWaterMarks
bool dequeue(U8 *buffer, FwSizeType capacity, FwSizeType &actualSize)
Dequeue a message (multi-consumer safe, non-blocking, O(1))
Status send(const U8 *buffer, FwSizeType size, FwQueuePriorityType priority, BlockingType blockType) override
send a message into the queue
PriorityMemQueue()
queue interface constructor - initializes handle
QueuePriorityConfig * priorityConfigs
FwEnumStoreType instanceId
FwSizeType m_numActivePriorities
FwSizeType getMessageHighWaterMark() const override
get maximum messages stored at any given time