749 using value_type = T;
753 typedef typename Traits::index_t index_t;
754 typedef typename Traits::size_t size_t;
756 static const size_t BLOCK_SIZE =
static_cast<size_t>(Traits::BLOCK_SIZE);
757 static const size_t EXPLICIT_BLOCK_EMPTY_COUNTER_THRESHOLD =
static_cast<size_t>(Traits::EXPLICIT_BLOCK_EMPTY_COUNTER_THRESHOLD);
758 static const size_t EXPLICIT_INITIAL_INDEX_SIZE =
static_cast<size_t>(Traits::EXPLICIT_INITIAL_INDEX_SIZE);
759 static const size_t IMPLICIT_INITIAL_INDEX_SIZE =
static_cast<size_t>(Traits::IMPLICIT_INITIAL_INDEX_SIZE);
760 static const size_t INITIAL_IMPLICIT_PRODUCER_HASH_SIZE =
static_cast<size_t>(Traits::INITIAL_IMPLICIT_PRODUCER_HASH_SIZE);
761 static const std::uint32_t EXPLICIT_CONSUMER_CONSUMPTION_QUOTA_BEFORE_ROTATE =
static_cast<std::uint32_t
>(Traits::EXPLICIT_CONSUMER_CONSUMPTION_QUOTA_BEFORE_ROTATE);
764#pragma warning(disable: 4307)
765#pragma warning(disable: 4309)
772 static_assert(!std::numeric_limits<size_t>::is_signed && std::is_integral<size_t>::value,
"Traits::size_t must be an unsigned integral type");
773 static_assert(!std::numeric_limits<index_t>::is_signed && std::is_integral<index_t>::value,
"Traits::index_t must be an unsigned integral type");
774 static_assert(
sizeof(index_t) >=
sizeof(
size_t),
"Traits::index_t must be at least as wide as Traits::size_t");
775 static_assert((BLOCK_SIZE > 1) && !(BLOCK_SIZE & (BLOCK_SIZE - 1)),
"Traits::BLOCK_SIZE must be a power of 2 (and at least 2)");
776 static_assert((EXPLICIT_BLOCK_EMPTY_COUNTER_THRESHOLD > 1) && !(EXPLICIT_BLOCK_EMPTY_COUNTER_THRESHOLD & (EXPLICIT_BLOCK_EMPTY_COUNTER_THRESHOLD - 1)),
"Traits::EXPLICIT_BLOCK_EMPTY_COUNTER_THRESHOLD must be a power of 2 (and greater than 1)");
777 static_assert((EXPLICIT_INITIAL_INDEX_SIZE > 1) && !(EXPLICIT_INITIAL_INDEX_SIZE & (EXPLICIT_INITIAL_INDEX_SIZE - 1)),
"Traits::EXPLICIT_INITIAL_INDEX_SIZE must be a power of 2 (and greater than 1)");
778 static_assert((IMPLICIT_INITIAL_INDEX_SIZE > 1) && !(IMPLICIT_INITIAL_INDEX_SIZE & (IMPLICIT_INITIAL_INDEX_SIZE - 1)),
"Traits::IMPLICIT_INITIAL_INDEX_SIZE must be a power of 2 (and greater than 1)");
779 static_assert((INITIAL_IMPLICIT_PRODUCER_HASH_SIZE == 0) || !(INITIAL_IMPLICIT_PRODUCER_HASH_SIZE & (INITIAL_IMPLICIT_PRODUCER_HASH_SIZE - 1)),
"Traits::INITIAL_IMPLICIT_PRODUCER_HASH_SIZE must be a power of 2");
780 static_assert(INITIAL_IMPLICIT_PRODUCER_HASH_SIZE == 0 || INITIAL_IMPLICIT_PRODUCER_HASH_SIZE >= 1,
"Traits::INITIAL_IMPLICIT_PRODUCER_HASH_SIZE must be at least 1 (or 0 to disable implicit enqueueing)");
794 : producerListTail(
nullptr),
796 initialBlockPoolIndex(0),
797 nextExplicitConsumerId(0),
798 globalExplicitConsumerOffset(0)
800 implicitProducerHashResizeInProgress.clear(std::memory_order_relaxed);
801 populate_initial_implicit_producer_hash();
802 populate_initial_block_list(capacity / BLOCK_SIZE + ((capacity & (BLOCK_SIZE - 1)) == 0 ? 0 : 1));
804#ifdef MOODYCAMEL_QUEUE_INTERNAL_DEBUG
818 : producerListTail(
nullptr),
820 initialBlockPoolIndex(0),
821 nextExplicitConsumerId(0),
822 globalExplicitConsumerOffset(0)
824 implicitProducerHashResizeInProgress.clear(std::memory_order_relaxed);
825 populate_initial_implicit_producer_hash();
827 populate_initial_block_list(blocks);
829#ifdef MOODYCAMEL_QUEUE_INTERNAL_DEBUG
841 auto ptr = producerListTail.load(std::memory_order_relaxed);
842 while (ptr !=
nullptr) {
843 auto next = ptr->next_prod();
844 if (ptr->token !=
nullptr) {
845 ptr->token->producer =
nullptr;
852 MOODYCAMEL_CONSTEXPR_IF (INITIAL_IMPLICIT_PRODUCER_HASH_SIZE != 0) {
853 auto hash = implicitProducerHash.load(std::memory_order_relaxed);
854 while (hash !=
nullptr) {
855 auto prev = hash->prev;
856 if (prev !=
nullptr) {
857 for (
size_t i = 0; i != hash->capacity; ++i) {
858 hash->entries[i].~ImplicitProducerKVP();
860 hash->~ImplicitProducerHash();
861 (Traits::free)(hash);
868 auto block = freeList.head_unsafe();
869 while (block !=
nullptr) {
870 auto next = block->freeListNext.load(std::memory_order_relaxed);
871 if (block->dynamicallyAllocated) {
878 destroy_array(initialBlockPool, initialBlockPoolSize);
892 : producerListTail(other.producerListTail.load(std::memory_order_relaxed)),
893 producerCount(other.producerCount.load(std::memory_order_relaxed)),
894 initialBlockPoolIndex(other.initialBlockPoolIndex.load(std::memory_order_relaxed)),
895 initialBlockPool(other.initialBlockPool),
896 initialBlockPoolSize(other.initialBlockPoolSize),
897 freeList(std::move(other.freeList)),
898 nextExplicitConsumerId(other.nextExplicitConsumerId.load(std::memory_order_relaxed)),
899 globalExplicitConsumerOffset(other.globalExplicitConsumerOffset.load(std::memory_order_relaxed))
902 implicitProducerHashResizeInProgress.clear(std::memory_order_relaxed);
903 populate_initial_implicit_producer_hash();
904 swap_implicit_producer_hashes(other);
906 other.producerListTail.store(
nullptr, std::memory_order_relaxed);
907 other.producerCount.store(0, std::memory_order_relaxed);
908 other.nextExplicitConsumerId.store(0, std::memory_order_relaxed);
909 other.globalExplicitConsumerOffset.store(0, std::memory_order_relaxed);
911#ifdef MOODYCAMEL_QUEUE_INTERNAL_DEBUG
912 explicitProducers.store(other.explicitProducers.load(std::memory_order_relaxed), std::memory_order_relaxed);
913 other.explicitProducers.store(
nullptr, std::memory_order_relaxed);
914 implicitProducers.store(other.implicitProducers.load(std::memory_order_relaxed), std::memory_order_relaxed);
915 other.implicitProducers.store(
nullptr, std::memory_order_relaxed);
918 other.initialBlockPoolIndex.store(0, std::memory_order_relaxed);
919 other.initialBlockPoolSize = 0;
920 other.initialBlockPool =
nullptr;
927 return swap_internal(other);
937 swap_internal(other);
943 if (
this == &other) {
947 details::swap_relaxed(producerListTail, other.producerListTail);
948 details::swap_relaxed(producerCount, other.producerCount);
949 details::swap_relaxed(initialBlockPoolIndex, other.initialBlockPoolIndex);
950 std::swap(initialBlockPool, other.initialBlockPool);
951 std::swap(initialBlockPoolSize, other.initialBlockPoolSize);
952 freeList.swap(other.freeList);
953 details::swap_relaxed(nextExplicitConsumerId, other.nextExplicitConsumerId);
954 details::swap_relaxed(globalExplicitConsumerOffset, other.globalExplicitConsumerOffset);
956 swap_implicit_producer_hashes(other);
959 other.reown_producers();
961#ifdef MOODYCAMEL_QUEUE_INTERNAL_DEBUG
975 inline bool enqueue(T
const&
item)
977 MOODYCAMEL_CONSTEXPR_IF (INITIAL_IMPLICIT_PRODUCER_HASH_SIZE == 0)
return false;
986 inline bool enqueue(T&&
item)
988 MOODYCAMEL_CONSTEXPR_IF (INITIAL_IMPLICIT_PRODUCER_HASH_SIZE == 0)
return false;
1016 template<
typename It>
1019 MOODYCAMEL_CONSTEXPR_IF (INITIAL_IMPLICIT_PRODUCER_HASH_SIZE == 0)
return false;
1029 template<
typename It>
1040 inline bool try_enqueue(T
const&
item)
1042 MOODYCAMEL_CONSTEXPR_IF (INITIAL_IMPLICIT_PRODUCER_HASH_SIZE == 0)
return false;
1051 inline bool try_enqueue(T&&
item)
1053 MOODYCAMEL_CONSTEXPR_IF (INITIAL_IMPLICIT_PRODUCER_HASH_SIZE == 0)
return false;
1080 template<
typename It>
1083 MOODYCAMEL_CONSTEXPR_IF (INITIAL_IMPLICIT_PRODUCER_HASH_SIZE == 0)
return false;
1092 template<
typename It>
1104 template<
typename U>
1105 bool try_dequeue(U&
item)
1110 ProducerBase*
best =
nullptr;
1112 for (
auto ptr = producerListTail.load(std::memory_order_acquire);
nonEmptyCount < 3 && ptr !=
nullptr; ptr = ptr->next_prod()) {
1113 auto size = ptr->size_approx();
1126 if ((details::likely)(
best->dequeue(
item))) {
1129 for (
auto ptr = producerListTail.load(std::memory_order_acquire); ptr !=
nullptr; ptr = ptr->next_prod()) {
1130 if (ptr !=
best && ptr->dequeue(
item)) {
1147 template<
typename U>
1148 bool try_dequeue_non_interleaved(U&
item)
1150 for (
auto ptr = producerListTail.load(std::memory_order_acquire); ptr !=
nullptr; ptr = ptr->next_prod()) {
1151 if (ptr->dequeue(
item)) {
1162 template<
typename U>
1171 if (token.desiredProducer ==
nullptr || token.lastKnownGlobalOffset != globalExplicitConsumerOffset.load(std::memory_order_relaxed)) {
1172 if (!update_current_producer_after_rotation(token)) {
1179 if (
static_cast<ProducerBase*
>(token.currentProducer)->dequeue(
item)) {
1180 if (++token.itemsConsumedFromCurrent == EXPLICIT_CONSUMER_CONSUMPTION_QUOTA_BEFORE_ROTATE) {
1181 globalExplicitConsumerOffset.fetch_add(1, std::memory_order_relaxed);
1186 auto tail = producerListTail.load(std::memory_order_acquire);
1187 auto ptr =
static_cast<ProducerBase*
>(token.currentProducer)->next_prod();
1188 if (ptr ==
nullptr) {
1191 while (ptr !=
static_cast<ProducerBase*
>(token.currentProducer)) {
1192 if (ptr->dequeue(
item)) {
1193 token.currentProducer = ptr;
1194 token.itemsConsumedFromCurrent = 1;
1197 ptr = ptr->next_prod();
1198 if (ptr ==
nullptr) {
1210 template<
typename It>
1214 for (
auto ptr = producerListTail.load(std::memory_order_acquire); ptr !=
nullptr; ptr = ptr->next_prod()) {
1215 count += ptr->dequeue_bulk(
itemFirst, max - count);
1228 template<
typename It>
1231 if (token.desiredProducer ==
nullptr || token.lastKnownGlobalOffset != globalExplicitConsumerOffset.load(std::memory_order_relaxed)) {
1232 if (!update_current_producer_after_rotation(token)) {
1237 size_t count =
static_cast<ProducerBase*
>(token.currentProducer)->dequeue_bulk(
itemFirst, max);
1239 if ((token.itemsConsumedFromCurrent +=
static_cast<std::uint32_t
>(max)) >= EXPLICIT_CONSUMER_CONSUMPTION_QUOTA_BEFORE_ROTATE) {
1240 globalExplicitConsumerOffset.fetch_add(1, std::memory_order_relaxed);
1244 token.itemsConsumedFromCurrent +=
static_cast<std::uint32_t
>(count);
1247 auto tail = producerListTail.load(std::memory_order_acquire);
1248 auto ptr =
static_cast<ProducerBase*
>(token.currentProducer)->next_prod();
1249 if (ptr ==
nullptr) {
1252 while (ptr !=
static_cast<ProducerBase*
>(token.currentProducer)) {
1256 token.currentProducer = ptr;
1257 token.itemsConsumedFromCurrent =
static_cast<std::uint32_t
>(
dequeued);
1263 ptr = ptr->next_prod();
1264 if (ptr ==
nullptr) {
1279 template<
typename U>
1282 return static_cast<ExplicitProducer*
>(producer.producer)->dequeue(
item);
1292 template<
typename It>
1295 return static_cast<ExplicitProducer*
>(producer.producer)->dequeue_bulk(
itemFirst, max);
1305 size_t size_approx()
const
1308 for (
auto ptr = producerListTail.load(std::memory_order_acquire); ptr !=
nullptr; ptr = ptr->next_prod()) {
1309 size += ptr->size_approx();
1318 static bool is_lock_free()
1333 struct ExplicitProducer;
1334 friend struct ExplicitProducer;
1335 struct ImplicitProducer;
1336 friend struct ImplicitProducer;
1337 friend class ConcurrentQueueTests;
1339 enum AllocationMode { CanAlloc, CannotAlloc };
1346 template<AllocationMode canAlloc,
typename U>
1349 return static_cast<ExplicitProducer*
>(token.producer)->ConcurrentQueue::ExplicitProducer::template
enqueue<canAlloc>(std::forward<U>(element));
1352 template<AllocationMode canAlloc,
typename U>
1353 inline bool inner_enqueue(U&& element)
1355 auto producer = get_or_add_implicit_producer();
1356 return producer ==
nullptr ? false : producer->ConcurrentQueue::ImplicitProducer::template
enqueue<canAlloc>(std::forward<U>(element));
1359 template<AllocationMode canAlloc,
typename It>
1365 template<AllocationMode canAlloc,
typename It>
1366 inline bool inner_enqueue_bulk(
It itemFirst,
size_t count)
1368 auto producer = get_or_add_implicit_producer();
1372 inline bool update_current_producer_after_rotation(
consumer_token_t& token)
1375 auto tail = producerListTail.load(std::memory_order_acquire);
1376 if (token.desiredProducer ==
nullptr && tail ==
nullptr) {
1379 auto prodCount = producerCount.load(std::memory_order_relaxed);
1380 auto globalOffset = globalExplicitConsumerOffset.load(std::memory_order_relaxed);
1381 if ((details::unlikely)(token.desiredProducer ==
nullptr)) {
1386 token.desiredProducer = tail;
1387 for (std::uint32_t i = 0; i != offset; ++i) {
1388 token.desiredProducer =
static_cast<ProducerBase*
>(token.desiredProducer)->next_prod();
1389 if (token.desiredProducer ==
nullptr) {
1390 token.desiredProducer = tail;
1395 std::uint32_t delta =
globalOffset - token.lastKnownGlobalOffset;
1399 for (std::uint32_t i = 0; i != delta; ++i) {
1400 token.desiredProducer =
static_cast<ProducerBase*
>(token.desiredProducer)->next_prod();
1401 if (token.desiredProducer ==
nullptr) {
1402 token.desiredProducer = tail;
1407 token.currentProducer = token.desiredProducer;
1408 token.itemsConsumedFromCurrent = 0;
1417 template <
typename N>
1420 FreeListNode() : freeListRefs(0), freeListNext(
nullptr) { }
1422 std::atomic<std::uint32_t> freeListRefs;
1423 std::atomic<N*> freeListNext;
1429 template<
typename N>
1432 FreeList() : freeListHead(
nullptr) { }
1433 FreeList(FreeList&& other) : freeListHead(other.freeListHead.load(std::memory_order_relaxed)) { other.freeListHead.store(
nullptr, std::memory_order_relaxed); }
1434 void swap(FreeList& other) { details::swap_relaxed(freeListHead, other.freeListHead); }
1436 FreeList(FreeList
const&) MOODYCAMEL_DELETE_FUNCTION;
1437 FreeList& operator=(FreeList
const&) MOODYCAMEL_DELETE_FUNCTION;
1439 inline void add(N* node)
1441#ifdef MCDBGQ_NOLOCKFREE_FREELIST
1442 debug::DebugLock lock(mutex);
1446 if (node->freeListRefs.fetch_add(SHOULD_BE_ON_FREELIST, std::memory_order_acq_rel) == 0) {
1449 add_knowing_refcount_is_zero(node);
1455#ifdef MCDBGQ_NOLOCKFREE_FREELIST
1456 debug::DebugLock lock(mutex);
1458 auto head = freeListHead.load(std::memory_order_acquire);
1459 while (head !=
nullptr) {
1461 auto refs = head->freeListRefs.load(std::memory_order_relaxed);
1462 if ((
refs & REFS_MASK) == 0 || !head->freeListRefs.compare_exchange_strong(
refs,
refs + 1, std::memory_order_acquire, std::memory_order_relaxed)) {
1463 head = freeListHead.load(std::memory_order_acquire);
1469 auto next = head->freeListNext.load(std::memory_order_relaxed);
1470 if (freeListHead.compare_exchange_strong(head, next, std::memory_order_acquire, std::memory_order_relaxed)) {
1473 assert((head->freeListRefs.load(std::memory_order_relaxed) & SHOULD_BE_ON_FREELIST) == 0);
1476 head->freeListRefs.fetch_sub(2, std::memory_order_release);
1483 refs =
prevHead->freeListRefs.fetch_sub(1, std::memory_order_acq_rel);
1484 if (
refs == SHOULD_BE_ON_FREELIST + 1) {
1485 add_knowing_refcount_is_zero(
prevHead);
1493 N* head_unsafe()
const {
return freeListHead.load(std::memory_order_relaxed); }
1496 inline void add_knowing_refcount_is_zero(N* node)
1506 auto head = freeListHead.load(std::memory_order_relaxed);
1508 node->freeListNext.store(head, std::memory_order_relaxed);
1509 node->freeListRefs.store(1, std::memory_order_release);
1510 if (!freeListHead.compare_exchange_strong(head, node, std::memory_order_release, std::memory_order_relaxed)) {
1512 if (node->freeListRefs.fetch_add(SHOULD_BE_ON_FREELIST - 1, std::memory_order_release) == 1) {
1522 std::atomic<N*> freeListHead;
1524 static const std::uint32_t REFS_MASK = 0x7FFFFFFF;
1525 static const std::uint32_t SHOULD_BE_ON_FREELIST = 0x80000000;
1527#ifdef MCDBGQ_NOLOCKFREE_FREELIST
1537 enum InnerQueueContext { implicit_context = 0, explicit_context = 1 };
1542 : next(
nullptr), elementsCompletelyDequeued(0), freeListRefs(0), freeListNext(
nullptr), shouldBeOnFreeList(
false), dynamicallyAllocated(
true)
1544#ifdef MCDBGQ_TRACKMEM
1549 template<InnerQueueContext context>
1550 inline bool is_empty()
const
1552 MOODYCAMEL_CONSTEXPR_IF (
context == explicit_context && BLOCK_SIZE <= EXPLICIT_BLOCK_EMPTY_COUNTER_THRESHOLD) {
1554 for (
size_t i = 0; i < BLOCK_SIZE; ++i) {
1555 if (!emptyFlags[i].load(std::memory_order_relaxed)) {
1561 std::atomic_thread_fence(std::memory_order_acquire);
1566 if (elementsCompletelyDequeued.load(std::memory_order_relaxed) == BLOCK_SIZE) {
1567 std::atomic_thread_fence(std::memory_order_acquire);
1570 assert(elementsCompletelyDequeued.load(std::memory_order_relaxed) <= BLOCK_SIZE);
1576 template<InnerQueueContext context>
1577 inline bool set_empty(MOODYCAMEL_MAYBE_UNUSED index_t i)
1579 MOODYCAMEL_CONSTEXPR_IF (
context == explicit_context && BLOCK_SIZE <= EXPLICIT_BLOCK_EMPTY_COUNTER_THRESHOLD) {
1581 assert(!emptyFlags[BLOCK_SIZE - 1 -
static_cast<size_t>(i &
static_cast<index_t
>(BLOCK_SIZE - 1))].load(std::memory_order_relaxed));
1582 emptyFlags[BLOCK_SIZE - 1 -
static_cast<size_t>(i &
static_cast<index_t
>(BLOCK_SIZE - 1))].store(
true, std::memory_order_release);
1587 auto prevVal = elementsCompletelyDequeued.fetch_add(1, std::memory_order_release);
1589 return prevVal == BLOCK_SIZE - 1;
1595 template<InnerQueueContext context>
1596 inline bool set_many_empty(MOODYCAMEL_MAYBE_UNUSED index_t i,
size_t count)
1598 MOODYCAMEL_CONSTEXPR_IF (
context == explicit_context && BLOCK_SIZE <= EXPLICIT_BLOCK_EMPTY_COUNTER_THRESHOLD) {
1600 std::atomic_thread_fence(std::memory_order_release);
1601 i = BLOCK_SIZE - 1 -
static_cast<size_t>(i &
static_cast<index_t
>(BLOCK_SIZE - 1)) - count + 1;
1602 for (
size_t j = 0;
j != count; ++
j) {
1603 assert(!emptyFlags[i +
j].load(std::memory_order_relaxed));
1604 emptyFlags[i +
j].store(
true, std::memory_order_relaxed);
1610 auto prevVal = elementsCompletelyDequeued.fetch_add(count, std::memory_order_release);
1612 return prevVal + count == BLOCK_SIZE;
1616 template<InnerQueueContext context>
1617 inline void set_all_empty()
1619 MOODYCAMEL_CONSTEXPR_IF (
context == explicit_context && BLOCK_SIZE <= EXPLICIT_BLOCK_EMPTY_COUNTER_THRESHOLD) {
1621 for (
size_t i = 0; i != BLOCK_SIZE; ++i) {
1622 emptyFlags[i].store(
true, std::memory_order_relaxed);
1627 elementsCompletelyDequeued.store(BLOCK_SIZE, std::memory_order_relaxed);
1631 template<InnerQueueContext context>
1632 inline void reset_empty()
1634 MOODYCAMEL_CONSTEXPR_IF (
context == explicit_context && BLOCK_SIZE <= EXPLICIT_BLOCK_EMPTY_COUNTER_THRESHOLD) {
1636 for (
size_t i = 0; i != BLOCK_SIZE; ++i) {
1637 emptyFlags[i].store(
false, std::memory_order_relaxed);
1642 elementsCompletelyDequeued.store(0, std::memory_order_relaxed);
1646 inline T* operator[](index_t idx) MOODYCAMEL_NOEXCEPT {
return static_cast<T*
>(
static_cast<void*
>(elements)) +
static_cast<size_t>(idx &
static_cast<index_t
>(BLOCK_SIZE - 1)); }
1647 inline T
const* operator[](index_t idx)
const MOODYCAMEL_NOEXCEPT {
return static_cast<T
const*
>(
static_cast<void const*
>(elements)) +
static_cast<size_t>(idx &
static_cast<index_t
>(BLOCK_SIZE - 1)); }
1650 static_assert(std::alignment_of<T>::value <=
sizeof(T),
"The queue does not support types with an alignment greater than their size at this time");
1651 MOODYCAMEL_ALIGNED_TYPE_LIKE(
char[
sizeof(T) * BLOCK_SIZE], T) elements;
1654 std::atomic<size_t> elementsCompletelyDequeued;
1655 std::atomic<bool> emptyFlags[BLOCK_SIZE <= EXPLICIT_BLOCK_EMPTY_COUNTER_THRESHOLD ? BLOCK_SIZE : 1];
1657 std::atomic<std::uint32_t> freeListRefs;
1658 std::atomic<Block*> freeListNext;
1659 std::atomic<bool> shouldBeOnFreeList;
1660 bool dynamicallyAllocated;
1662#ifdef MCDBGQ_TRACKMEM
1666 static_assert(std::alignment_of<Block>::value >= std::alignment_of<T>::value,
"Internal error: Blocks must be at least as aligned as the type they are wrapping");
1669#ifdef MCDBGQ_TRACKMEM
1684 dequeueOptimisticCount(0),
1685 dequeueOvercommit(0),
1692 virtual ~ProducerBase() { }
1694 template<
typename U>
1695 inline bool dequeue(U& element)
1698 return static_cast<ExplicitProducer*
>(
this)->dequeue(element);
1701 return static_cast<ImplicitProducer*
>(
this)->dequeue(element);
1705 template<
typename It>
1706 inline size_t dequeue_bulk(
It&
itemFirst,
size_t max)
1709 return static_cast<ExplicitProducer*
>(
this)->dequeue_bulk(
itemFirst, max);
1712 return static_cast<ImplicitProducer*
>(
this)->dequeue_bulk(
itemFirst, max);
1716 inline ProducerBase* next_prod()
const {
return static_cast<ProducerBase*
>(next); }
1718 inline size_t size_approx()
const
1720 auto tail = tailIndex.load(std::memory_order_relaxed);
1721 auto head = headIndex.load(std::memory_order_relaxed);
1722 return details::circular_less_than(head, tail) ?
static_cast<size_t>(tail - head) : 0;
1725 inline index_t getTail()
const {
return tailIndex.load(std::memory_order_relaxed); }
1727 std::atomic<index_t> tailIndex;
1728 std::atomic<index_t> headIndex;
1730 std::atomic<index_t> dequeueOptimisticCount;
1731 std::atomic<index_t> dequeueOvercommit;
1740#ifdef MCDBGQ_TRACKMEM
1750 struct ExplicitProducer :
public ProducerBase
1753 ProducerBase(parent_,
true),
1754 blockIndex(
nullptr),
1755 pr_blockIndexSlotsUsed(0),
1756 pr_blockIndexSize(EXPLICIT_INITIAL_INDEX_SIZE >> 1),
1757 pr_blockIndexFront(0),
1758 pr_blockIndexEntries(
nullptr),
1759 pr_blockIndexRaw(
nullptr)
1761 size_t poolBasedIndexSize = details::ceil_to_pow_2(parent_->initialBlockPoolSize) >> 1;
1774 if (this->tailBlock !=
nullptr) {
1777 if ((this->headIndex.load(std::memory_order_relaxed) &
static_cast<index_t
>(BLOCK_SIZE - 1)) != 0) {
1780 size_t i = (pr_blockIndexFront - pr_blockIndexSlotsUsed) & (pr_blockIndexSize - 1);
1781 while (details::circular_less_than<index_t>(pr_blockIndexEntries[i].base + BLOCK_SIZE, this->headIndex.load(std::memory_order_relaxed))) {
1782 i = (i + 1) & (pr_blockIndexSize - 1);
1784 assert(details::circular_less_than<index_t>(pr_blockIndexEntries[i].base, this->headIndex.load(std::memory_order_relaxed)));
1789 auto block = this->tailBlock;
1791 block = block->next;
1798 i =
static_cast<size_t>(this->headIndex.load(std::memory_order_relaxed) &
static_cast<index_t
>(BLOCK_SIZE - 1));
1802 auto lastValidIndex = (this->tailIndex.load(std::memory_order_relaxed) &
static_cast<index_t
>(BLOCK_SIZE - 1)) == 0 ? BLOCK_SIZE :
static_cast<size_t>(this->tailIndex.load(std::memory_order_relaxed) &
static_cast<index_t
>(BLOCK_SIZE - 1));
1803 while (i != BLOCK_SIZE && (block != this->tailBlock || i !=
lastValidIndex)) {
1804 (*block)[i++]->~T();
1806 }
while (block != this->tailBlock);
1810 if (this->tailBlock !=
nullptr) {
1811 auto block = this->tailBlock;
1814 if (block->dynamicallyAllocated) {
1818 this->parent->add_block_to_free_list(block);
1821 }
while (block != this->tailBlock);
1825 auto header =
static_cast<BlockIndexHeader*
>(pr_blockIndexRaw);
1826 while (header !=
nullptr) {
1827 auto prev =
static_cast<BlockIndexHeader*
>(header->prev);
1828 header->~BlockIndexHeader();
1829 (Traits::free)(header);
1834 template<AllocationMode allocMode,
typename U>
1835 inline bool enqueue(U&& element)
1837 index_t
currentTailIndex = this->tailIndex.load(std::memory_order_relaxed);
1845 this->tailBlock = this->tailBlock->next;
1858 auto head = this->headIndex.load(std::memory_order_relaxed);
1860 if (!details::circular_less_than<index_t>(head,
currentTailIndex + BLOCK_SIZE)
1868 if (pr_blockIndexRaw ==
nullptr || pr_blockIndexSlotsUsed == pr_blockIndexSize) {
1873 MOODYCAMEL_CONSTEXPR_IF (
allocMode == CannotAlloc) {
1876 else if (!new_block_index(pr_blockIndexSlotsUsed)) {
1886#ifdef MCDBGQ_TRACKMEM
1890 if (this->tailBlock ==
nullptr) {
1894 newBlock->next = this->tailBlock->next;
1898 ++pr_blockIndexSlotsUsed;
1901 MOODYCAMEL_CONSTEXPR_IF (!MOODYCAMEL_NOEXCEPT_CTOR(T, U,
new (
static_cast<T*
>(
nullptr)) T(std::forward<U>(element)))) {
1907 MOODYCAMEL_CATCH (...) {
1921 auto& entry = blockIndex.load(std::memory_order_relaxed)->entries[pr_blockIndexFront];
1923 entry.block = this->tailBlock;
1924 blockIndex.load(std::memory_order_relaxed)->front.store(pr_blockIndexFront, std::memory_order_release);
1925 pr_blockIndexFront = (pr_blockIndexFront + 1) & (pr_blockIndexSize - 1);
1927 MOODYCAMEL_CONSTEXPR_IF (!MOODYCAMEL_NOEXCEPT_CTOR(T, U,
new (
static_cast<T*
>(
nullptr)) T(std::forward<U>(element)))) {
1928 this->tailIndex.store(
newTailIndex, std::memory_order_release);
1936 this->tailIndex.store(
newTailIndex, std::memory_order_release);
1940 template<
typename U>
1941 bool dequeue(U& element)
1943 auto tail = this->tailIndex.load(std::memory_order_relaxed);
1944 auto overcommit = this->dequeueOvercommit.load(std::memory_order_relaxed);
1945 if (details::circular_less_than<index_t>(this->dequeueOptimisticCount.load(std::memory_order_relaxed) -
overcommit, tail)) {
1962 std::atomic_thread_fence(std::memory_order_acquire);
1965 auto myDequeueCount = this->dequeueOptimisticCount.fetch_add(1, std::memory_order_relaxed);
1977 tail = this->tailIndex.load(std::memory_order_acquire);
1989 auto index = this->headIndex.fetch_add(1, std::memory_order_acq_rel);
2006 auto&
el = *((*block)[index]);
2007 if (!MOODYCAMEL_NOEXCEPT_ASSIGN(T, T&&, element = std::move(
el))) {
2016 (*block)[index]->~T();
2019 }
guard = { block, index };
2021 element = std::move(
el);
2024 element = std::move(
el);
2033 this->dequeueOvercommit.fetch_add(1, std::memory_order_release);
2040 template<AllocationMode allocMode,
typename It>
2041 bool MOODYCAMEL_NO_TSAN enqueue_bulk(
It itemFirst,
size_t count)
2046 index_t
startTailIndex = this->tailIndex.load(std::memory_order_relaxed);
2062 this->tailBlock = this->tailBlock->next;
2065 auto& entry = blockIndex.load(std::memory_order_relaxed)->entries[pr_blockIndexFront];
2067 entry.block = this->tailBlock;
2068 pr_blockIndexFront = (pr_blockIndexFront + 1) & (pr_blockIndexSize - 1);
2076 auto head = this->headIndex.load(std::memory_order_relaxed);
2079 if (pr_blockIndexRaw ==
nullptr || pr_blockIndexSlotsUsed == pr_blockIndexSize ||
full) {
2080 MOODYCAMEL_CONSTEXPR_IF (
allocMode == CannotAlloc) {
2110#ifdef MCDBGQ_TRACKMEM
2114 if (this->tailBlock ==
nullptr) {
2118 newBlock->next = this->tailBlock->next;
2124 ++pr_blockIndexSlotsUsed;
2126 auto& entry = blockIndex.load(std::memory_order_relaxed)->entries[pr_blockIndexFront];
2128 entry.block = this->tailBlock;
2129 pr_blockIndexFront = (pr_blockIndexFront + 1) & (pr_blockIndexSize - 1);
2137 if (block == this->tailBlock) {
2140 block = block->next;
2143 MOODYCAMEL_CONSTEXPR_IF (MOODYCAMEL_NOEXCEPT_CTOR(T,
decltype(*
itemFirst),
new (
static_cast<T*
>(
nullptr)) T(details::deref_noexcept(
itemFirst)))) {
2144 blockIndex.load(std::memory_order_relaxed)->front.store((pr_blockIndexFront - 1) & (pr_blockIndexSize - 1), std::memory_order_release);
2162 MOODYCAMEL_CONSTEXPR_IF (MOODYCAMEL_NOEXCEPT_CTOR(T,
decltype(*
itemFirst),
new (
static_cast<T*
>(
nullptr)) T(details::deref_noexcept(
itemFirst)))) {
2182 MOODYCAMEL_CATCH (...) {
2195 if ((
startTailIndex &
static_cast<index_t
>(BLOCK_SIZE - 1)) == 0) {
2210 block = block->next;
2221 this->tailBlock = this->tailBlock->next;
2224 MOODYCAMEL_CONSTEXPR_IF (!MOODYCAMEL_NOEXCEPT_CTOR(T,
decltype(*
itemFirst),
new (
static_cast<T*
>(
nullptr)) T(details::deref_noexcept(
itemFirst)))) {
2226 blockIndex.load(std::memory_order_relaxed)->front.store((pr_blockIndexFront - 1) & (pr_blockIndexSize - 1), std::memory_order_release);
2229 this->tailIndex.store(
newTailIndex, std::memory_order_release);
2233 template<
typename It>
2236 auto tail = this->tailIndex.load(std::memory_order_relaxed);
2237 auto overcommit = this->dequeueOvercommit.load(std::memory_order_relaxed);
2238 auto desiredCount =
static_cast<size_t>(tail - (this->dequeueOptimisticCount.load(std::memory_order_relaxed) -
overcommit));
2239 if (details::circular_less_than<size_t>(0,
desiredCount)) {
2241 std::atomic_thread_fence(std::memory_order_acquire);
2245 tail = this->tailIndex.load(std::memory_order_acquire);
2247 if (details::circular_less_than<size_t>(0,
actualCount)) {
2255 auto firstIndex = this->headIndex.fetch_add(
actualCount, std::memory_order_acq_rel);
2267 auto index = firstIndex;
2270 index_t endIndex = (index &
~static_cast<index_t>(BLOCK_SIZE - 1)) +
static_cast<index_t
>(BLOCK_SIZE);
2271 endIndex = details::circular_less_than<index_t>(firstIndex +
static_cast<index_t
>(
actualCount), endIndex) ? firstIndex +
static_cast<index_t
>(
actualCount) : endIndex;
2273 if (MOODYCAMEL_NOEXCEPT_ASSIGN(T, T&&, details::deref_noexcept(
itemFirst) = std::move((*(*block)[index])))) {
2274 while (index != endIndex) {
2275 auto&
el = *((*block)[index]);
2283 while (index != endIndex) {
2284 auto&
el = *((*block)[index]);
2291 MOODYCAMEL_CATCH (...) {
2297 while (index != endIndex) {
2298 (*block)[index++]->~T();
2305 endIndex = details::circular_less_than<index_t>(firstIndex +
static_cast<index_t
>(
actualCount), endIndex) ? firstIndex +
static_cast<index_t
>(
actualCount) : endIndex;
2319 this->dequeueOvercommit.fetch_add(
desiredCount, std::memory_order_release);
2327 struct BlockIndexEntry
2333 struct BlockIndexHeader
2336 std::atomic<size_t> front;
2337 BlockIndexEntry* entries;
2347 pr_blockIndexSize <<= 1;
2348 auto newRawPtr =
static_cast<char*
>((Traits::malloc)(
sizeof(BlockIndexHeader) + std::alignment_of<BlockIndexEntry>::value - 1 +
sizeof(BlockIndexEntry) * pr_blockIndexSize));
2350 pr_blockIndexSize >>= 1;
2358 if (pr_blockIndexSlotsUsed != 0) {
2363 }
while (i != pr_blockIndexFront);
2367 auto header =
new (
newRawPtr) BlockIndexHeader;
2368 header->size = pr_blockIndexSize;
2371 header->prev = pr_blockIndexRaw;
2373 pr_blockIndexFront =
j;
2376 blockIndex.store(header, std::memory_order_release);
2382 std::atomic<BlockIndexHeader*> blockIndex;
2385 size_t pr_blockIndexSlotsUsed;
2386 size_t pr_blockIndexSize;
2387 size_t pr_blockIndexFront;
2388 BlockIndexEntry* pr_blockIndexEntries;
2389 void* pr_blockIndexRaw;
2391#ifdef MOODYCAMEL_QUEUE_INTERNAL_DEBUG
2397#ifdef MCDBGQ_TRACKMEM
2407 struct ImplicitProducer :
public ProducerBase
2410 ProducerBase(parent_,
false),
2411 nextBlockIndexCapacity(IMPLICIT_INITIAL_INDEX_SIZE),
2424#ifdef MOODYCAMEL_CPP11_THREAD_LOCAL_SUPPORTED
2426 if (!this->inactive.load(std::memory_order_relaxed)) {
2432 auto tail = this->tailIndex.load(std::memory_order_relaxed);
2433 auto index = this->headIndex.load(std::memory_order_relaxed);
2434 Block* block =
nullptr;
2435 assert(index == tail || details::circular_less_than(index, tail));
2437 while (index != tail) {
2438 if ((index &
static_cast<index_t
>(BLOCK_SIZE - 1)) == 0 || block ==
nullptr) {
2439 if (block !=
nullptr) {
2441 this->parent->add_block_to_free_list(block);
2444 block = get_block_index_entry_for_index(index)->value.load(std::memory_order_relaxed);
2447 ((*block)[index])->~T();
2453 if (this->tailBlock !=
nullptr && (
forceFreeLastBlock || (tail &
static_cast<index_t
>(BLOCK_SIZE - 1)) != 0)) {
2454 this->parent->add_block_to_free_list(this->tailBlock);
2472 template<AllocationMode allocMode,
typename U>
2473 inline bool enqueue(U&& element)
2475 index_t
currentTailIndex = this->tailIndex.load(std::memory_order_relaxed);
2479 auto head = this->headIndex.load(std::memory_order_relaxed);
2484#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODBLOCKINDEX
2485 debug::DebugLock lock(mutex);
2496 rewind_block_index_tail();
2497 idxEntry->value.store(
nullptr, std::memory_order_relaxed);
2500#ifdef MCDBGQ_TRACKMEM
2505 MOODYCAMEL_CONSTEXPR_IF (!MOODYCAMEL_NOEXCEPT_CTOR(T, U,
new (
static_cast<T*
>(
nullptr)) T(std::forward<U>(element)))) {
2510 MOODYCAMEL_CATCH (...) {
2511 rewind_block_index_tail();
2512 idxEntry->value.store(
nullptr, std::memory_order_relaxed);
2513 this->parent->add_block_to_free_list(
newBlock);
2523 MOODYCAMEL_CONSTEXPR_IF (!MOODYCAMEL_NOEXCEPT_CTOR(T, U,
new (
static_cast<T*
>(
nullptr)) T(std::forward<U>(element)))) {
2524 this->tailIndex.store(
newTailIndex, std::memory_order_release);
2532 this->tailIndex.store(
newTailIndex, std::memory_order_release);
2536 template<
typename U>
2537 bool dequeue(U& element)
2540 index_t tail = this->tailIndex.load(std::memory_order_relaxed);
2541 index_t
overcommit = this->dequeueOvercommit.load(std::memory_order_relaxed);
2542 if (details::circular_less_than<index_t>(this->dequeueOptimisticCount.load(std::memory_order_relaxed) -
overcommit, tail)) {
2543 std::atomic_thread_fence(std::memory_order_acquire);
2545 index_t
myDequeueCount = this->dequeueOptimisticCount.fetch_add(1, std::memory_order_relaxed);
2546 tail = this->tailIndex.load(std::memory_order_acquire);
2548 index_t index = this->headIndex.fetch_add(1, std::memory_order_acq_rel);
2551 auto entry = get_block_index_entry_for_index(index);
2554 auto block = entry->value.load(std::memory_order_relaxed);
2555 auto&
el = *((*block)[index]);
2557 if (!MOODYCAMEL_NOEXCEPT_ASSIGN(T, T&&, element = std::move(
el))) {
2558#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODBLOCKINDEX
2561 debug::DebugLock lock(producer->mutex);
2566 BlockIndexEntry* entry;
2571 (*block)[index]->~T();
2573 entry->value.store(
nullptr, std::memory_order_relaxed);
2574 parent->add_block_to_free_list(block);
2577 }
guard = { block, index, entry, this->parent };
2579 element = std::move(
el);
2582 element = std::move(
el);
2587#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODBLOCKINDEX
2588 debug::DebugLock lock(mutex);
2591 entry->value.store(
nullptr, std::memory_order_relaxed);
2593 this->parent->add_block_to_free_list(block);
2600 this->dequeueOvercommit.fetch_add(1, std::memory_order_release);
2608#pragma warning(push)
2609#pragma warning(disable: 4706)
2611 template<AllocationMode allocMode,
typename It>
2623 index_t
startTailIndex = this->tailIndex.load(std::memory_order_relaxed);
2632#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODBLOCKINDEX
2633 debug::DebugLock lock(mutex);
2640 BlockIndexEntry*
idxEntry =
nullptr;
2643 auto head = this->headIndex.load(std::memory_order_relaxed);
2651 rewind_block_index_tail();
2652 idxEntry->value.store(
nullptr, std::memory_order_relaxed);
2658 idxEntry->value.store(
nullptr, std::memory_order_relaxed);
2659 rewind_block_index_tail();
2667#ifdef MCDBGQ_TRACKMEM
2679 assert(this->tailBlock !=
nullptr);
2701 MOODYCAMEL_CONSTEXPR_IF (MOODYCAMEL_NOEXCEPT_CTOR(T,
decltype(*
itemFirst),
new (
static_cast<T*
>(
nullptr)) T(details::deref_noexcept(
itemFirst)))) {
2714 MOODYCAMEL_CATCH (...) {
2720 if ((
startTailIndex &
static_cast<index_t
>(BLOCK_SIZE - 1)) == 0) {
2735 block = block->next;
2743 idxEntry->value.store(
nullptr, std::memory_order_relaxed);
2744 rewind_block_index_tail();
2756 this->tailBlock = this->tailBlock->next;
2758 this->tailIndex.store(
newTailIndex, std::memory_order_release);
2765 template<
typename It>
2768 auto tail = this->tailIndex.load(std::memory_order_relaxed);
2769 auto overcommit = this->dequeueOvercommit.load(std::memory_order_relaxed);
2770 auto desiredCount =
static_cast<size_t>(tail - (this->dequeueOptimisticCount.load(std::memory_order_relaxed) -
overcommit));
2771 if (details::circular_less_than<size_t>(0,
desiredCount)) {
2773 std::atomic_thread_fence(std::memory_order_acquire);
2777 tail = this->tailIndex.load(std::memory_order_acquire);
2779 if (details::circular_less_than<size_t>(0,
actualCount)) {
2787 auto firstIndex = this->headIndex.fetch_add(
actualCount, std::memory_order_acq_rel);
2790 auto index = firstIndex;
2795 index_t endIndex = (index &
~static_cast<index_t>(BLOCK_SIZE - 1)) +
static_cast<index_t
>(BLOCK_SIZE);
2796 endIndex = details::circular_less_than<index_t>(firstIndex +
static_cast<index_t
>(
actualCount), endIndex) ? firstIndex +
static_cast<index_t
>(
actualCount) : endIndex;
2799 auto block = entry->value.load(std::memory_order_relaxed);
2800 if (MOODYCAMEL_NOEXCEPT_ASSIGN(T, T&&, details::deref_noexcept(
itemFirst) = std::move((*(*block)[index])))) {
2801 while (index != endIndex) {
2802 auto&
el = *((*block)[index]);
2810 while (index != endIndex) {
2811 auto&
el = *((*block)[index]);
2818 MOODYCAMEL_CATCH (...) {
2821 block = entry->value.load(std::memory_order_relaxed);
2822 while (index != endIndex) {
2823 (*block)[index++]->~T();
2827#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODBLOCKINDEX
2828 debug::DebugLock lock(mutex);
2830 entry->value.store(
nullptr, std::memory_order_relaxed);
2831 this->parent->add_block_to_free_list(block);
2837 endIndex = details::circular_less_than<index_t>(firstIndex +
static_cast<index_t
>(
actualCount), endIndex) ? firstIndex +
static_cast<index_t
>(
actualCount) : endIndex;
2845#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODBLOCKINDEX
2846 debug::DebugLock lock(mutex);
2850 entry->value.store(
nullptr, std::memory_order_relaxed);
2852 this->parent->add_block_to_free_list(block);
2860 this->dequeueOvercommit.fetch_add(
desiredCount, std::memory_order_release);
2869 static const index_t INVALID_BLOCK_BASE = 1;
2871 struct BlockIndexEntry
2873 std::atomic<index_t> key;
2874 std::atomic<Block*> value;
2877 struct BlockIndexHeader
2880 std::atomic<size_t> tail;
2881 BlockIndexEntry* entries;
2882 BlockIndexEntry** index;
2883 BlockIndexHeader* prev;
2886 template<AllocationMode allocMode>
2895 if (
idxEntry->key.load(std::memory_order_relaxed) == INVALID_BLOCK_BASE ||
2896 idxEntry->value.load(std::memory_order_relaxed) ==
nullptr) {
2904 MOODYCAMEL_CONSTEXPR_IF (
allocMode == CannotAlloc) {
2907 else if (!new_block_index()) {
2913 assert(
idxEntry->key.load(std::memory_order_relaxed) == INVALID_BLOCK_BASE);
2919 inline void rewind_block_index_tail()
2925 inline BlockIndexEntry* get_block_index_entry_for_index(index_t index)
const
2932 inline size_t get_block_index_index_for_index(index_t index, BlockIndexHeader*&
localBlockIndex)
const
2934#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODBLOCKINDEX
2935 debug::DebugLock lock(mutex);
2944 auto offset =
static_cast<size_t>(
static_cast<typename std::make_signed<index_t>::type
>(index -
tailBase) / BLOCK_SIZE);
2950 bool new_block_index()
2952 auto prev = blockIndex.load(std::memory_order_relaxed);
2953 size_t prevCapacity = prev ==
nullptr ? 0 : prev->capacity;
2955 auto raw =
static_cast<char*
>((Traits::malloc)(
2956 sizeof(BlockIndexHeader) +
2957 std::alignment_of<BlockIndexEntry>::value - 1 +
sizeof(BlockIndexEntry) *
entryCount +
2958 std::alignment_of<BlockIndexEntry*>::value - 1 +
sizeof(BlockIndexEntry*) * nextBlockIndexCapacity));
2959 if (raw ==
nullptr) {
2963 auto header =
new (raw) BlockIndexHeader;
2964 auto entries =
reinterpret_cast<BlockIndexEntry*
>(details::align_for<BlockIndexEntry>(raw +
sizeof(BlockIndexHeader)));
2965 auto index =
reinterpret_cast<BlockIndexEntry**
>(details::align_for<BlockIndexEntry*>(
reinterpret_cast<char*
>(entries) +
sizeof(BlockIndexEntry) *
entryCount));
2966 if (prev !=
nullptr) {
2967 auto prevTail = prev->tail.load(std::memory_order_relaxed);
2972 index[i++] = prev->index[
prevPos];
2977 new (entries + i) BlockIndexEntry;
2978 entries[i].key.store(INVALID_BLOCK_BASE, std::memory_order_relaxed);
2981 header->prev = prev;
2982 header->entries = entries;
2983 header->index = index;
2984 header->capacity = nextBlockIndexCapacity;
2985 header->tail.store((
prevCapacity - 1) & (nextBlockIndexCapacity - 1), std::memory_order_relaxed);
2987 blockIndex.store(header, std::memory_order_release);
2989 nextBlockIndexCapacity <<= 1;
2995 size_t nextBlockIndexCapacity;
2996 std::atomic<BlockIndexHeader*> blockIndex;
2998#ifdef MOODYCAMEL_CPP11_THREAD_LOCAL_SUPPORTED
3004#ifdef MOODYCAMEL_QUEUE_INTERNAL_DEBUG
3010#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODBLOCKINDEX
3013#ifdef MCDBGQ_TRACKMEM
3023 void populate_initial_block_list(
size_t blockCount)
3025 initialBlockPoolSize = blockCount;
3026 if (initialBlockPoolSize == 0) {
3027 initialBlockPool =
nullptr;
3032 if (initialBlockPool ==
nullptr) {
3033 initialBlockPoolSize = 0;
3035 for (
size_t i = 0; i < initialBlockPoolSize; ++i) {
3036 initialBlockPool[i].dynamicallyAllocated =
false;
3040 inline Block* try_get_block_from_initial_pool()
3042 if (initialBlockPoolIndex.load(std::memory_order_relaxed) >= initialBlockPoolSize) {
3046 auto index = initialBlockPoolIndex.fetch_add(1, std::memory_order_relaxed);
3048 return index < initialBlockPoolSize ? (initialBlockPool + index) :
nullptr;
3051 inline void add_block_to_free_list(Block* block)
3053#ifdef MCDBGQ_TRACKMEM
3054 block->owner =
nullptr;
3056 freeList.add(block);
3059 inline void add_blocks_to_free_list(Block* block)
3061 while (block !=
nullptr) {
3062 auto next = block->next;
3063 add_block_to_free_list(block);
3068 inline Block* try_get_block_from_free_list()
3070 return freeList.try_get();
3074 template<AllocationMode canAlloc>
3075 Block* requisition_block()
3077 auto block = try_get_block_from_initial_pool();
3078 if (block !=
nullptr) {
3082 block = try_get_block_from_free_list();
3083 if (block !=
nullptr) {
3087 MOODYCAMEL_CONSTEXPR_IF (
canAlloc == CanAlloc) {
3096#ifdef MCDBGQ_TRACKMEM
3119 stats.elementsEnqueued = q->size_approx();
3121 auto block = q->freeList.head_unsafe();
3122 while (block !=
nullptr) {
3123 ++stats.allocatedBlocks;
3125 block = block->freeListNext.load(std::memory_order_relaxed);
3128 for (
auto ptr = q->producerListTail.load(std::memory_order_acquire); ptr !=
nullptr; ptr = ptr->next_prod()) {
3129 bool implicit =
dynamic_cast<ImplicitProducer*
>(ptr) !=
nullptr;
3130 stats.implicitProducers +=
implicit ? 1 : 0;
3131 stats.explicitProducers +=
implicit ? 0 : 1;
3134 auto prod =
static_cast<ImplicitProducer*
>(ptr);
3135 stats.queueClassBytes +=
sizeof(ImplicitProducer);
3136 auto head = prod->headIndex.load(std::memory_order_relaxed);
3137 auto tail = prod->tailIndex.load(std::memory_order_relaxed);
3138 auto hash = prod->blockIndex.load(std::memory_order_relaxed);
3139 if (hash !=
nullptr) {
3140 for (
size_t i = 0; i != hash->capacity; ++i) {
3141 if (hash->index[i]->key.load(std::memory_order_relaxed) != ImplicitProducer::INVALID_BLOCK_BASE && hash->index[i]->value.load(std::memory_order_relaxed) !=
nullptr) {
3142 ++stats.allocatedBlocks;
3143 ++stats.ownedBlocksImplicit;
3146 stats.implicitBlockIndexBytes += hash->capacity *
sizeof(
typename ImplicitProducer::BlockIndexEntry);
3147 for (; hash !=
nullptr; hash = hash->prev) {
3148 stats.implicitBlockIndexBytes +=
sizeof(
typename ImplicitProducer::BlockIndexHeader) + hash->capacity *
sizeof(
typename ImplicitProducer::BlockIndexEntry*);
3151 for (; details::circular_less_than<index_t>(head, tail); head += BLOCK_SIZE) {
3157 auto prod =
static_cast<ExplicitProducer*
>(ptr);
3158 stats.queueClassBytes +=
sizeof(ExplicitProducer);
3159 auto tailBlock = prod->tailBlock;
3161 if (tailBlock !=
nullptr) {
3162 auto block = tailBlock;
3164 ++stats.allocatedBlocks;
3169 ++stats.ownedBlocksExplicit;
3170 block = block->next;
3171 }
while (block != tailBlock);
3173 auto index = prod->blockIndex.load(std::memory_order_relaxed);
3174 while (index !=
nullptr) {
3175 stats.explicitBlockIndexBytes +=
sizeof(
typename ExplicitProducer::BlockIndexHeader) + index->size *
sizeof(
typename ExplicitProducer::BlockIndexEntry);
3176 index =
static_cast<typename ExplicitProducer::BlockIndexHeader*
>(index->prev);
3181 auto freeOnInitialPool = q->initialBlockPoolIndex.load(std::memory_order_relaxed) >= q->initialBlockPoolSize ? 0 : q->initialBlockPoolSize - q->initialBlockPoolIndex.load(std::memory_order_relaxed);
3185 stats.blockClassBytes =
sizeof(Block) * stats.allocatedBlocks;
3195 return MemStats::getFor(
this);
3206 ProducerBase* recycle_or_create_producer(
bool isExplicit)
3209 return recycle_or_create_producer(isExplicit,
recycled);
3212 ProducerBase* recycle_or_create_producer(
bool isExplicit,
bool&
recycled)
3214#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODHASH
3218 for (
auto ptr = producerListTail.load(std::memory_order_acquire); ptr !=
nullptr; ptr = ptr->next_prod()) {
3219 if (ptr->inactive.load(std::memory_order_relaxed) && ptr->isExplicit == isExplicit) {
3220 bool expected =
true;
3221 if (ptr->inactive.compare_exchange_strong(expected,
false, std::memory_order_acquire, std::memory_order_relaxed)) {
3233 ProducerBase* add_producer(ProducerBase* producer)
3236 if (producer ==
nullptr) {
3240 producerCount.fetch_add(1, std::memory_order_relaxed);
3243 auto prevTail = producerListTail.load(std::memory_order_relaxed);
3246 }
while (!producerListTail.compare_exchange_weak(
prevTail, producer, std::memory_order_release, std::memory_order_relaxed));
3248#ifdef MOODYCAMEL_QUEUE_INTERNAL_DEBUG
3249 if (producer->isExplicit) {
3253 }
while (!
explicitProducers.compare_exchange_weak(
prevTailExplicit,
static_cast<ExplicitProducer*
>(producer), std::memory_order_release, std::memory_order_relaxed));
3259 }
while (!
implicitProducers.compare_exchange_weak(
prevTailImplicit,
static_cast<ImplicitProducer*
>(producer), std::memory_order_release, std::memory_order_relaxed));
3266 void reown_producers()
3271 for (
auto ptr = producerListTail.load(std::memory_order_relaxed); ptr !=
nullptr; ptr = ptr->next_prod()) {
3281 struct ImplicitProducerKVP
3283 std::atomic<details::thread_id_t> key;
3284 ImplicitProducer* value;
3286 ImplicitProducerKVP() : value(
nullptr) { }
3288 ImplicitProducerKVP(ImplicitProducerKVP&& other) MOODYCAMEL_NOEXCEPT
3290 key.store(other.key.load(std::memory_order_relaxed), std::memory_order_relaxed);
3291 value = other.value;
3294 inline ImplicitProducerKVP& operator=(ImplicitProducerKVP&& other) MOODYCAMEL_NOEXCEPT
3300 inline void swap(ImplicitProducerKVP& other) MOODYCAMEL_NOEXCEPT
3302 if (
this != &other) {
3303 details::swap_relaxed(key, other.key);
3304 std::swap(value, other.value);
3309 template<
typename XT,
typename XTraits>
3312 struct ImplicitProducerHash
3315 ImplicitProducerKVP* entries;
3316 ImplicitProducerHash* prev;
3319 inline void populate_initial_implicit_producer_hash()
3321 MOODYCAMEL_CONSTEXPR_IF (INITIAL_IMPLICIT_PRODUCER_HASH_SIZE == 0) {
3325 implicitProducerHashCount.store(0, std::memory_order_relaxed);
3326 auto hash = &initialImplicitProducerHash;
3327 hash->capacity = INITIAL_IMPLICIT_PRODUCER_HASH_SIZE;
3328 hash->entries = &initialImplicitProducerHashEntries[0];
3329 for (
size_t i = 0; i != INITIAL_IMPLICIT_PRODUCER_HASH_SIZE; ++i) {
3330 initialImplicitProducerHashEntries[i].key.store(details::invalid_thread_id, std::memory_order_relaxed);
3332 hash->prev =
nullptr;
3333 implicitProducerHash.store(hash, std::memory_order_relaxed);
3339 MOODYCAMEL_CONSTEXPR_IF (INITIAL_IMPLICIT_PRODUCER_HASH_SIZE == 0) {
3344 initialImplicitProducerHashEntries.swap(other.initialImplicitProducerHashEntries);
3345 initialImplicitProducerHash.entries = &initialImplicitProducerHashEntries[0];
3346 other.initialImplicitProducerHash.entries = &other.initialImplicitProducerHashEntries[0];
3348 details::swap_relaxed(implicitProducerHashCount, other.implicitProducerHashCount);
3350 details::swap_relaxed(implicitProducerHash, other.implicitProducerHash);
3351 if (implicitProducerHash.load(std::memory_order_relaxed) == &other.initialImplicitProducerHash) {
3352 implicitProducerHash.store(&initialImplicitProducerHash, std::memory_order_relaxed);
3355 ImplicitProducerHash* hash;
3356 for (hash = implicitProducerHash.load(std::memory_order_relaxed); hash->prev != &other.initialImplicitProducerHash; hash = hash->prev) {
3359 hash->prev = &initialImplicitProducerHash;
3361 if (other.implicitProducerHash.load(std::memory_order_relaxed) == &initialImplicitProducerHash) {
3362 other.implicitProducerHash.store(&other.initialImplicitProducerHash, std::memory_order_relaxed);
3365 ImplicitProducerHash* hash;
3366 for (hash = other.implicitProducerHash.load(std::memory_order_relaxed); hash->prev != &initialImplicitProducerHash; hash = hash->prev) {
3369 hash->prev = &other.initialImplicitProducerHash;
3375 ImplicitProducer* get_or_add_implicit_producer()
3387#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODHASH
3391 auto id = details::thread_id();
3392 auto hashedId = details::hash_thread_id(
id);
3394 auto mainHash = implicitProducerHash.load(std::memory_order_acquire);
3396 for (
auto hash =
mainHash; hash !=
nullptr; hash = hash->prev) {
3400 index &= hash->capacity - 1;
3402 auto probedKey = hash->entries[index].key.load(std::memory_order_relaxed);
3409 auto value = hash->entries[index].value;
3415 auto empty = details::invalid_thread_id;
3416#ifdef MOODYCAMEL_CPP11_THREAD_LOCAL_SUPPORTED
3417 auto reusable = details::invalid_thread_id2;
3418 if ((
probedKey == empty &&
mainHash->entries[index].key.compare_exchange_strong(empty,
id, std::memory_order_relaxed, std::memory_order_relaxed)) ||
3421 if ((
probedKey == empty &&
mainHash->entries[index].key.compare_exchange_strong(empty,
id, std::memory_order_relaxed, std::memory_order_relaxed))) {
3423 mainHash->entries[index].value = value;
3432 if (
probedKey == details::invalid_thread_id) {
3440 auto newCount = 1 + implicitProducerHashCount.fetch_add(1, std::memory_order_relaxed);
3443 if (
newCount >= (
mainHash->capacity >> 1) && !implicitProducerHashResizeInProgress.test_and_set(std::memory_order_acquire)) {
3448 mainHash = implicitProducerHash.load(std::memory_order_acquire);
3454 auto raw =
static_cast<char*
>((Traits::malloc)(
sizeof(ImplicitProducerHash) + std::alignment_of<ImplicitProducerKVP>::value - 1 +
sizeof(ImplicitProducerKVP) *
newCapacity));
3455 if (raw ==
nullptr) {
3457 implicitProducerHashCount.fetch_sub(1, std::memory_order_relaxed);
3458 implicitProducerHashResizeInProgress.clear(std::memory_order_relaxed);
3462 auto newHash =
new (raw) ImplicitProducerHash;
3464 newHash->entries =
reinterpret_cast<ImplicitProducerKVP*
>(details::align_for<ImplicitProducerKVP>(raw +
sizeof(ImplicitProducerHash)));
3466 new (
newHash->entries + i) ImplicitProducerKVP;
3467 newHash->entries[i].key.store(details::invalid_thread_id, std::memory_order_relaxed);
3470 implicitProducerHash.store(
newHash, std::memory_order_release);
3471 implicitProducerHashResizeInProgress.clear(std::memory_order_release);
3475 implicitProducerHashResizeInProgress.clear(std::memory_order_release);
3484 auto producer =
static_cast<ImplicitProducer*
>(recycle_or_create_producer(
false,
recycled));
3485 if (producer ==
nullptr) {
3486 implicitProducerHashCount.fetch_sub(1, std::memory_order_relaxed);
3490 implicitProducerHashCount.fetch_sub(1, std::memory_order_relaxed);
3493#ifdef MOODYCAMEL_CPP11_THREAD_LOCAL_SUPPORTED
3494 producer->threadExitListener.callback = &ConcurrentQueue::implicit_producer_thread_exited_callback;
3495 producer->threadExitListener.userData = producer;
3496 details::ThreadExitNotifier::subscribe(&producer->threadExitListener);
3504 auto empty = details::invalid_thread_id;
3505#ifdef MOODYCAMEL_CPP11_THREAD_LOCAL_SUPPORTED
3506 auto reusable = details::invalid_thread_id2;
3507 if ((
probedKey == empty &&
mainHash->entries[index].key.compare_exchange_strong(empty,
id, std::memory_order_relaxed, std::memory_order_relaxed)) ||
3510 if ((
probedKey == empty &&
mainHash->entries[index].key.compare_exchange_strong(empty,
id, std::memory_order_relaxed, std::memory_order_relaxed))) {
3512 mainHash->entries[index].value = producer;
3523 mainHash = implicitProducerHash.load(std::memory_order_acquire);
3527#ifdef MOODYCAMEL_CPP11_THREAD_LOCAL_SUPPORTED
3531 details::ThreadExitNotifier::unsubscribe(&producer->threadExitListener);
3534#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODHASH
3537 auto hash = implicitProducerHash.load(std::memory_order_acquire);
3539 auto id = details::thread_id();
3540 auto hashedId = details::hash_thread_id(
id);
3545 for (; hash !=
nullptr; hash = hash->prev) {
3548 index &= hash->capacity - 1;
3549 probedKey = hash->entries[index].key.load(std::memory_order_relaxed);
3551 hash->entries[index].key.store(details::invalid_thread_id2, std::memory_order_release);
3555 }
while (
probedKey != details::invalid_thread_id);
3559 producer->inactive.store(
true, std::memory_order_release);
3564 auto producer =
static_cast<ImplicitProducer*
>(userData);
3565 auto queue = producer->parent;
3566 queue->implicit_producer_thread_exited(producer);
3574 template<
typename TAlign>
3575 static inline void* aligned_malloc(
size_t size)
3577 MOODYCAMEL_CONSTEXPR_IF (std::alignment_of<TAlign>::value <= std::alignment_of<details::max_align_t>::value)
3578 return (Traits::malloc)(size);
3580 size_t alignment = std::alignment_of<TAlign>::value;
3581 void* raw = (Traits::malloc)(size + alignment - 1 +
sizeof(
void*));
3584 char* ptr = details::align_for<TAlign>(
reinterpret_cast<char*
>(raw) +
sizeof(
void*));
3585 *(
reinterpret_cast<void**
>(ptr) - 1) = raw;
3590 template<
typename TAlign>
3591 static inline void aligned_free(
void* ptr)
3593 MOODYCAMEL_CONSTEXPR_IF (std::alignment_of<TAlign>::value <= std::alignment_of<details::max_align_t>::value)
3594 return (Traits::free)(ptr);
3596 (Traits::free)(ptr ? *(
reinterpret_cast<void**
>(ptr) - 1) :
nullptr);
3599 template<
typename U>
3600 static inline U* create_array(
size_t count)
3607 for (
size_t i = 0; i != count; ++i)
3612 template<
typename U>
3613 static inline void destroy_array(U* p,
size_t count)
3617 for (
size_t i = count; i != 0; )
3623 template<
typename U>
3624 static inline U* create()
3627 return p !=
nullptr ?
new (p) U :
nullptr;
3630 template<
typename U,
typename A1>
3631 static inline U* create(
A1&&
a1)
3634 return p !=
nullptr ?
new (p) U(std::forward<A1>(
a1)) :
nullptr;
3637 template<
typename U>
3638 static inline void destroy(U* p)
3646 std::atomic<ProducerBase*> producerListTail;
3647 std::atomic<std::uint32_t> producerCount;
3649 std::atomic<size_t> initialBlockPoolIndex;
3650 Block* initialBlockPool;
3651 size_t initialBlockPoolSize;
3653#ifndef MCDBGQ_USEDEBUGFREELIST
3656 debug::DebugFreeList<Block> freeList;
3659 std::atomic<ImplicitProducerHash*> implicitProducerHash;
3660 std::atomic<size_t> implicitProducerHashCount;
3661 ImplicitProducerHash initialImplicitProducerHash;
3662 std::array<ImplicitProducerKVP, INITIAL_IMPLICIT_PRODUCER_HASH_SIZE> initialImplicitProducerHashEntries;
3663 std::atomic_flag implicitProducerHashResizeInProgress;
3665 std::atomic<std::uint32_t> nextExplicitConsumerId;
3666 std::atomic<std::uint32_t> globalExplicitConsumerOffset;
3668#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODHASH
3672#ifdef MOODYCAMEL_QUEUE_INTERNAL_DEBUG