RavEngine
Loading...
Searching...
No Matches
concurrentqueue.h
1// Provides a C++11 implementation of a multi-producer, multi-consumer lock-free queue.
2// An overview, including benchmark results, is provided here:
3// http://moodycamel.com/blog/2014/a-fast-general-purpose-lock-free-queue-for-c++
4// The full design is also described in excruciating detail at:
5// http://moodycamel.com/blog/2014/detailed-design-of-a-lock-free-queue
6
7// Simplified BSD license:
8// Copyright (c) 2013-2020, Cameron Desrochers.
9// All rights reserved.
10//
11// Redistribution and use in source and binary forms, with or without modification,
12// are permitted provided that the following conditions are met:
13//
14// - Redistributions of source code must retain the above copyright notice, this list of
15// conditions and the following disclaimer.
16// - Redistributions in binary form must reproduce the above copyright notice, this list of
17// conditions and the following disclaimer in the documentation and/or other materials
18// provided with the distribution.
19//
20// THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS IS" AND ANY
21// EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED WARRANTIES OF
22// MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL
23// THE COPYRIGHT HOLDER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL,
24// SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT
25// OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION)
26// HOWEVER CAUSED AND ON ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR
27// TORT (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS SOFTWARE,
28// EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
29
30// Also dual-licensed under the Boost Software License (see LICENSE.md)
31
32#pragma once
33
34#if defined(__GNUC__)
35// Disable -Wconversion warnings (spuriously triggered when Traits::size_t and
36// Traits::index_t are set to < 32 bits, causing integer promotion, causing warnings
37// upon assigning any computed values)
38#pragma GCC diagnostic push
39#pragma GCC diagnostic ignored "-Wconversion"
40
41#ifdef MCDBGQ_USE_RELACY
42#pragma GCC diagnostic ignored "-Wint-to-pointer-cast"
43#endif
44#endif
45
46#if defined(_MSC_VER) && (!defined(_HAS_CXX17) || !_HAS_CXX17)
47// VS2019 with /W4 warns about constant conditional expressions but unless /std=c++17 or higher
48// does not support `if constexpr`, so we have no choice but to simply disable the warning
49#pragma warning(push)
50#pragma warning(disable: 4127) // conditional expression is constant
51#endif
52
53#if defined(__APPLE__)
54#include "TargetConditionals.h"
55#endif
56
57#ifdef MCDBGQ_USE_RELACY
58#include "relacy/relacy_std.hpp"
59#include "relacy_shims.h"
60// We only use malloc/free anyway, and the delete macro messes up `= delete` method declarations.
61// We'll override the default trait malloc ourselves without a macro.
62#undef new
63#undef delete
64#undef malloc
65#undef free
66#else
67#include <atomic> // Requires C++11. Sorry VS2010.
68#include <cassert>
69#endif
70#include <cstddef> // for max_align_t
71#include <cstdint>
72#include <cstdlib>
73#include <type_traits>
74#include <algorithm>
75#include <utility>
76#include <limits>
77#include <climits> // for CHAR_BIT
78#include <array>
79#include <thread> // partly for __WINPTHREADS_VERSION if on MinGW-w64 w/ POSIX threading
80
81// Platform-specific definitions of a numeric thread ID type and an invalid value
82namespace moodycamel { namespace details {
83 template<typename thread_id_t> struct thread_id_converter {
84 typedef thread_id_t thread_id_numeric_size_t;
85 typedef thread_id_t thread_id_hash_t;
86 static thread_id_hash_t prehash(thread_id_t const& x) { return x; }
87 };
88} }
89#if defined(MCDBGQ_USE_RELACY)
90namespace moodycamel { namespace details {
91 typedef std::uint32_t thread_id_t;
92 static const thread_id_t invalid_thread_id = 0xFFFFFFFFU;
93 static const thread_id_t invalid_thread_id2 = 0xFFFFFFFEU;
94 static inline thread_id_t thread_id() { return rl::thread_index(); }
95} }
96#elif defined(_WIN32) || defined(__WINDOWS__) || defined(__WIN32__)
97// No sense pulling in windows.h in a header, we'll manually declare the function
98// we use and rely on backwards-compatibility for this not to break
99extern "C" __declspec(dllimport) unsigned long __stdcall GetCurrentThreadId(void);
100namespace moodycamel { namespace details {
101 static_assert(sizeof(unsigned long) == sizeof(std::uint32_t), "Expected size of unsigned long to be 32 bits on Windows");
102 typedef std::uint32_t thread_id_t;
103 static const thread_id_t invalid_thread_id = 0; // See http://blogs.msdn.com/b/oldnewthing/archive/2004/02/23/78395.aspx
104 static const thread_id_t invalid_thread_id2 = 0xFFFFFFFFU; // Not technically guaranteed to be invalid, but is never used in practice. Note that all Win32 thread IDs are presently multiples of 4.
105 static inline thread_id_t thread_id() { return static_cast<thread_id_t>(::GetCurrentThreadId()); }
106} }
107#elif defined(__arm__) || defined(_M_ARM) || defined(__aarch64__) || (defined(__APPLE__) && TARGET_OS_IPHONE)
108namespace moodycamel { namespace details {
109 static_assert(sizeof(std::thread::id) == 4 || sizeof(std::thread::id) == 8, "std::thread::id is expected to be either 4 or 8 bytes");
110
111 typedef std::thread::id thread_id_t;
112 static const thread_id_t invalid_thread_id; // Default ctor creates invalid ID
113
114 // Note we don't define a invalid_thread_id2 since std::thread::id doesn't have one; it's
115 // only used if MOODYCAMEL_CPP11_THREAD_LOCAL_SUPPORTED is defined anyway, which it won't
116 // be.
117 static inline thread_id_t thread_id() { return std::this_thread::get_id(); }
118
119 template<std::size_t> struct thread_id_size { };
120 template<> struct thread_id_size<4> { typedef std::uint32_t numeric_t; };
121 template<> struct thread_id_size<8> { typedef std::uint64_t numeric_t; };
122
123 template<> struct thread_id_converter<thread_id_t> {
124 typedef thread_id_size<sizeof(thread_id_t)>::numeric_t thread_id_numeric_size_t;
125#ifndef __APPLE__
126 typedef std::size_t thread_id_hash_t;
127#else
128 typedef thread_id_numeric_size_t thread_id_hash_t;
129#endif
130
131 static thread_id_hash_t prehash(thread_id_t const& x)
132 {
133#ifndef __APPLE__
134 return std::hash<std::thread::id>()(x);
135#else
136 return *reinterpret_cast<thread_id_hash_t const*>(&x);
137#endif
138 }
139 };
140} }
141#else
142// Use a nice trick from this answer: http://stackoverflow.com/a/8438730/21475
143// In order to get a numeric thread ID in a platform-independent way, we use a thread-local
144// static variable's address as a thread identifier :-)
145#if defined(__GNUC__) || defined(__INTEL_COMPILER)
146#define MOODYCAMEL_THREADLOCAL __thread
147#elif defined(_MSC_VER)
148#define MOODYCAMEL_THREADLOCAL __declspec(thread)
149#else
150// Assume C++11 compliant compiler
151#define MOODYCAMEL_THREADLOCAL thread_local
152#endif
153namespace moodycamel { namespace details {
154 typedef std::uintptr_t thread_id_t;
155 static const thread_id_t invalid_thread_id = 0; // Address can't be nullptr
156 static const thread_id_t invalid_thread_id2 = 1; // Member accesses off a null pointer are also generally invalid. Plus it's not aligned.
157 inline thread_id_t thread_id() { static MOODYCAMEL_THREADLOCAL int x; return reinterpret_cast<thread_id_t>(&x); }
158} }
159#endif
160
161// Constexpr if
162#ifndef MOODYCAMEL_CONSTEXPR_IF
163#if (defined(_MSC_VER) && defined(_HAS_CXX17) && _HAS_CXX17) || __cplusplus > 201402L
164#define MOODYCAMEL_CONSTEXPR_IF if constexpr
165#define MOODYCAMEL_MAYBE_UNUSED [[maybe_unused]]
166#else
167#define MOODYCAMEL_CONSTEXPR_IF if
168#define MOODYCAMEL_MAYBE_UNUSED
169#endif
170#endif
171
172// Exceptions
173#ifndef MOODYCAMEL_EXCEPTIONS_ENABLED
174#if (defined(_MSC_VER) && defined(_CPPUNWIND)) || (defined(__GNUC__) && defined(__EXCEPTIONS)) || (!defined(_MSC_VER) && !defined(__GNUC__))
175#define MOODYCAMEL_EXCEPTIONS_ENABLED
176#endif
177#endif
178#ifdef MOODYCAMEL_EXCEPTIONS_ENABLED
179#define MOODYCAMEL_TRY try
180#define MOODYCAMEL_CATCH(...) catch(__VA_ARGS__)
181#define MOODYCAMEL_RETHROW throw
182#define MOODYCAMEL_THROW(expr) throw (expr)
183#else
184#define MOODYCAMEL_TRY MOODYCAMEL_CONSTEXPR_IF (true)
185#define MOODYCAMEL_CATCH(...) else MOODYCAMEL_CONSTEXPR_IF (false)
186#define MOODYCAMEL_RETHROW
187#define MOODYCAMEL_THROW(expr)
188#endif
189
190#ifndef MOODYCAMEL_NOEXCEPT
191#if !defined(MOODYCAMEL_EXCEPTIONS_ENABLED)
192#define MOODYCAMEL_NOEXCEPT
193#define MOODYCAMEL_NOEXCEPT_CTOR(type, valueType, expr) true
194#define MOODYCAMEL_NOEXCEPT_ASSIGN(type, valueType, expr) true
195#elif defined(_MSC_VER) && defined(_NOEXCEPT) && _MSC_VER < 1800
196// VS2012's std::is_nothrow_[move_]constructible is broken and returns true when it shouldn't :-(
197// We have to assume *all* non-trivial constructors may throw on VS2012!
198#define MOODYCAMEL_NOEXCEPT _NOEXCEPT
199#define MOODYCAMEL_NOEXCEPT_CTOR(type, valueType, expr) (std::is_rvalue_reference<valueType>::value && std::is_move_constructible<type>::value ? std::is_trivially_move_constructible<type>::value : std::is_trivially_copy_constructible<type>::value)
200#define MOODYCAMEL_NOEXCEPT_ASSIGN(type, valueType, expr) ((std::is_rvalue_reference<valueType>::value && std::is_move_assignable<type>::value ? std::is_trivially_move_assignable<type>::value || std::is_nothrow_move_assignable<type>::value : std::is_trivially_copy_assignable<type>::value || std::is_nothrow_copy_assignable<type>::value) && MOODYCAMEL_NOEXCEPT_CTOR(type, valueType, expr))
201#elif defined(_MSC_VER) && defined(_NOEXCEPT) && _MSC_VER < 1900
202#define MOODYCAMEL_NOEXCEPT _NOEXCEPT
203#define MOODYCAMEL_NOEXCEPT_CTOR(type, valueType, expr) (std::is_rvalue_reference<valueType>::value && std::is_move_constructible<type>::value ? std::is_trivially_move_constructible<type>::value || std::is_nothrow_move_constructible<type>::value : std::is_trivially_copy_constructible<type>::value || std::is_nothrow_copy_constructible<type>::value)
204#define MOODYCAMEL_NOEXCEPT_ASSIGN(type, valueType, expr) ((std::is_rvalue_reference<valueType>::value && std::is_move_assignable<type>::value ? std::is_trivially_move_assignable<type>::value || std::is_nothrow_move_assignable<type>::value : std::is_trivially_copy_assignable<type>::value || std::is_nothrow_copy_assignable<type>::value) && MOODYCAMEL_NOEXCEPT_CTOR(type, valueType, expr))
205#else
206#define MOODYCAMEL_NOEXCEPT noexcept
207#define MOODYCAMEL_NOEXCEPT_CTOR(type, valueType, expr) noexcept(expr)
208#define MOODYCAMEL_NOEXCEPT_ASSIGN(type, valueType, expr) noexcept(expr)
209#endif
210#endif
211
212#ifndef MOODYCAMEL_CPP11_THREAD_LOCAL_SUPPORTED
213#ifdef MCDBGQ_USE_RELACY
214#define MOODYCAMEL_CPP11_THREAD_LOCAL_SUPPORTED
215#else
216// VS2013 doesn't support `thread_local`, and MinGW-w64 w/ POSIX threading has a crippling bug: http://sourceforge.net/p/mingw-w64/bugs/445
217// g++ <=4.7 doesn't support thread_local either.
218// Finally, iOS/ARM doesn't have support for it either, and g++/ARM allows it to compile but it's unconfirmed to actually work
219#if (!defined(_MSC_VER) || _MSC_VER >= 1900) && (!defined(__MINGW32__) && !defined(__MINGW64__) || !defined(__WINPTHREADS_VERSION)) && (!defined(__GNUC__) || __GNUC__ > 4 || (__GNUC__ == 4 && __GNUC_MINOR__ >= 8)) && (!defined(__APPLE__) || !TARGET_OS_IPHONE) && !defined(__arm__) && !defined(_M_ARM) && !defined(__aarch64__)
220// Assume `thread_local` is fully supported in all other C++11 compilers/platforms
221//#define MOODYCAMEL_CPP11_THREAD_LOCAL_SUPPORTED // always disabled for now since several users report having problems with it on
222#endif
223#endif
224#endif
225
226// VS2012 doesn't support deleted functions.
227// In this case, we declare the function normally but don't define it. A link error will be generated if the function is called.
228#ifndef MOODYCAMEL_DELETE_FUNCTION
229#if defined(_MSC_VER) && _MSC_VER < 1800
230#define MOODYCAMEL_DELETE_FUNCTION
231#else
232#define MOODYCAMEL_DELETE_FUNCTION = delete
233#endif
234#endif
235
236namespace moodycamel { namespace details {
237#ifndef MOODYCAMEL_ALIGNAS
238// VS2013 doesn't support alignas or alignof, and align() requires a constant literal
239#if defined(_MSC_VER) && _MSC_VER <= 1800
240#define MOODYCAMEL_ALIGNAS(alignment) __declspec(align(alignment))
241#define MOODYCAMEL_ALIGNOF(obj) __alignof(obj)
242#define MOODYCAMEL_ALIGNED_TYPE_LIKE(T, obj) typename details::Vs2013Aligned<std::alignment_of<obj>::value, T>::type
243 template<int Align, typename T> struct Vs2013Aligned { }; // default, unsupported alignment
244 template<typename T> struct Vs2013Aligned<1, T> { typedef __declspec(align(1)) T type; };
245 template<typename T> struct Vs2013Aligned<2, T> { typedef __declspec(align(2)) T type; };
246 template<typename T> struct Vs2013Aligned<4, T> { typedef __declspec(align(4)) T type; };
247 template<typename T> struct Vs2013Aligned<8, T> { typedef __declspec(align(8)) T type; };
248 template<typename T> struct Vs2013Aligned<16, T> { typedef __declspec(align(16)) T type; };
249 template<typename T> struct Vs2013Aligned<32, T> { typedef __declspec(align(32)) T type; };
250 template<typename T> struct Vs2013Aligned<64, T> { typedef __declspec(align(64)) T type; };
251 template<typename T> struct Vs2013Aligned<128, T> { typedef __declspec(align(128)) T type; };
252 template<typename T> struct Vs2013Aligned<256, T> { typedef __declspec(align(256)) T type; };
253#else
254 template<typename T> struct identity { typedef T type; };
255#define MOODYCAMEL_ALIGNAS(alignment) alignas(alignment)
256#define MOODYCAMEL_ALIGNOF(obj) alignof(obj)
257#define MOODYCAMEL_ALIGNED_TYPE_LIKE(T, obj) alignas(alignof(obj)) typename details::identity<T>::type
258#endif
259#endif
260} }
261
262
263// TSAN can false report races in lock-free code. To enable TSAN to be used from projects that use this one,
264// we can apply per-function compile-time suppression.
265// See https://clang.llvm.org/docs/ThreadSanitizer.html#has-feature-thread-sanitizer
266#define MOODYCAMEL_NO_TSAN
267#if defined(__has_feature)
268 #if __has_feature(thread_sanitizer)
269 #undef MOODYCAMEL_NO_TSAN
270 #define MOODYCAMEL_NO_TSAN __attribute__((no_sanitize("thread")))
271 #endif // TSAN
272#endif // TSAN
273
274// Compiler-specific likely/unlikely hints
275namespace moodycamel { namespace details {
276#if defined(__GNUC__)
277 static inline bool (likely)(bool x) { return __builtin_expect((x), true); }
278 static inline bool (unlikely)(bool x) { return __builtin_expect((x), false); }
279#else
280 static inline bool (likely)(bool x) { return x; }
281 static inline bool (unlikely)(bool x) { return x; }
282#endif
283} }
284
285#ifdef MOODYCAMEL_QUEUE_INTERNAL_DEBUG
286#include "internal/concurrentqueue_internal_debug.h"
287#endif
288
289namespace moodycamel {
290namespace details {
291 template<typename T>
293 static_assert(std::is_integral<T>::value, "const_numeric_max can only be used with integers");
294 static const T value = std::numeric_limits<T>::is_signed
295 ? (static_cast<T>(1) << (sizeof(T) * CHAR_BIT - 1)) - static_cast<T>(1)
296 : static_cast<T>(-1);
297 };
298
299#if defined(__GLIBCXX__)
300 typedef ::max_align_t std_max_align_t; // libstdc++ forgot to add it to std:: for a while
301#else
302 typedef std::max_align_t std_max_align_t; // Others (e.g. MSVC) insist it can *only* be accessed via std::
303#endif
304
305 // Some platforms have incorrectly set max_align_t to a type with <8 bytes alignment even while supporting
306 // 8-byte aligned scalar values (*cough* 32-bit iOS). Work around this with our own union. See issue #64.
307 typedef union {
308 std_max_align_t x;
309 long long y;
310 void* z;
311 } max_align_t;
312}
313
314// Default traits for the ConcurrentQueue. To change some of the
315// traits without re-implementing all of them, inherit from this
316// struct and shadow the declarations you wish to be different;
317// since the traits are used as a template type parameter, the
318// shadowed declarations will be used where defined, and the defaults
319// otherwise.
321{
322 // General-purpose size type. std::size_t is strongly recommended.
323 typedef std::size_t size_t;
324
325 // The type used for the enqueue and dequeue indices. Must be at least as
326 // large as size_t. Should be significantly larger than the number of elements
327 // you expect to hold at once, especially if you have a high turnover rate;
328 // for example, on 32-bit x86, if you expect to have over a hundred million
329 // elements or pump several million elements through your queue in a very
330 // short space of time, using a 32-bit type *may* trigger a race condition.
331 // A 64-bit int type is recommended in that case, and in practice will
332 // prevent a race condition no matter the usage of the queue. Note that
333 // whether the queue is lock-free with a 64-int type depends on the whether
334 // std::atomic<std::uint64_t> is lock-free, which is platform-specific.
335 typedef std::size_t index_t;
336
337 // Internally, all elements are enqueued and dequeued from multi-element
338 // blocks; this is the smallest controllable unit. If you expect few elements
339 // but many producers, a smaller block size should be favoured. For few producers
340 // and/or many elements, a larger block size is preferred. A sane default
341 // is provided. Must be a power of 2.
342 static const size_t BLOCK_SIZE = 32;
343
344 // For explicit producers (i.e. when using a producer token), the block is
345 // checked for being empty by iterating through a list of flags, one per element.
346 // For large block sizes, this is too inefficient, and switching to an atomic
347 // counter-based approach is faster. The switch is made for block sizes strictly
348 // larger than this threshold.
349 static const size_t EXPLICIT_BLOCK_EMPTY_COUNTER_THRESHOLD = 32;
350
351 // How many full blocks can be expected for a single explicit producer? This should
352 // reflect that number's maximum for optimal performance. Must be a power of 2.
353 static const size_t EXPLICIT_INITIAL_INDEX_SIZE = 32;
354
355 // How many full blocks can be expected for a single implicit producer? This should
356 // reflect that number's maximum for optimal performance. Must be a power of 2.
357 static const size_t IMPLICIT_INITIAL_INDEX_SIZE = 32;
358
359 // The initial size of the hash table mapping thread IDs to implicit producers.
360 // Note that the hash is resized every time it becomes half full.
361 // Must be a power of two, and either 0 or at least 1. If 0, implicit production
362 // (using the enqueue methods without an explicit producer token) is disabled.
363 static const size_t INITIAL_IMPLICIT_PRODUCER_HASH_SIZE = 32;
364
365 // Controls the number of items that an explicit consumer (i.e. one with a token)
366 // must consume before it causes all consumers to rotate and move on to the next
367 // internal queue.
368 static const std::uint32_t EXPLICIT_CONSUMER_CONSUMPTION_QUOTA_BEFORE_ROTATE = 256;
369
370 // The maximum number of elements (inclusive) that can be enqueued to a sub-queue.
371 // Enqueue operations that would cause this limit to be surpassed will fail. Note
372 // that this limit is enforced at the block level (for performance reasons), i.e.
373 // it's rounded up to the nearest block size.
374 static const size_t MAX_SUBQUEUE_SIZE = details::const_numeric_max<size_t>::value;
375
376 // The number of times to spin before sleeping when waiting on a semaphore.
377 // Recommended values are on the order of 1000-10000 unless the number of
378 // consumer threads exceeds the number of idle cores (in which case try 0-100).
379 // Only affects instances of the BlockingConcurrentQueue.
380 static const int MAX_SEMA_SPINS = 10000;
381
382
383#ifndef MCDBGQ_USE_RELACY
384 // Memory allocation can be customized if needed.
385 // malloc should return nullptr on failure, and handle alignment like std::malloc.
386#if defined(malloc) || defined(free)
387 // Gah, this is 2015, stop defining macros that break standard code already!
388 // Work around malloc/free being special macros:
389 static inline void* WORKAROUND_malloc(size_t size) { return malloc(size); }
390 static inline void WORKAROUND_free(void* ptr) { return free(ptr); }
391 static inline void* (malloc)(size_t size) { return WORKAROUND_malloc(size); }
392 static inline void (free)(void* ptr) { return WORKAROUND_free(ptr); }
393#else
394 static inline void* malloc(size_t size) { return std::malloc(size); }
395 static inline void free(void* ptr) { return std::free(ptr); }
396#endif
397#else
398 // Debug versions when running under the Relacy race detector (ignore
399 // these in user code)
400 static inline void* malloc(size_t size) { return rl::rl_malloc(size, $); }
401 static inline void free(void* ptr) { return rl::rl_free(ptr, $); }
402#endif
403};
404
405
406// When producing or consuming many elements, the most efficient way is to:
407// 1) Use one of the bulk-operation methods of the queue with a token
408// 2) Failing that, use the bulk-operation methods without a token
409// 3) Failing that, create a token and use that with the single-item methods
410// 4) Failing that, use the single-parameter methods of the queue
411// Having said that, don't create tokens willy-nilly -- ideally there should be
412// a maximum of one token per thread (of each kind).
413struct ProducerToken;
414struct ConsumerToken;
415
416template<typename T, typename Traits> class ConcurrentQueue;
417template<typename T, typename Traits> class BlockingConcurrentQueue;
418class ConcurrentQueueTests;
419
420
421namespace details
422{
424 {
426 std::atomic<bool> inactive;
427 ProducerToken* token;
428
430 : next(nullptr), inactive(false), token(nullptr)
431 {
432 }
433 };
434
435 template<bool use32> struct _hash_32_or_64 {
436 static inline std::uint32_t hash(std::uint32_t h)
437 {
438 // MurmurHash3 finalizer -- see https://code.google.com/p/smhasher/source/browse/trunk/MurmurHash3.cpp
439 // Since the thread ID is already unique, all we really want to do is propagate that
440 // uniqueness evenly across all the bits, so that we can use a subset of the bits while
441 // reducing collisions significantly
442 h ^= h >> 16;
443 h *= 0x85ebca6b;
444 h ^= h >> 13;
445 h *= 0xc2b2ae35;
446 return h ^ (h >> 16);
447 }
448 };
449 template<> struct _hash_32_or_64<1> {
450 static inline std::uint64_t hash(std::uint64_t h)
451 {
452 h ^= h >> 33;
453 h *= 0xff51afd7ed558ccd;
454 h ^= h >> 33;
455 h *= 0xc4ceb9fe1a85ec53;
456 return h ^ (h >> 33);
457 }
458 };
459 template<std::size_t size> struct hash_32_or_64 : public _hash_32_or_64<(size > 4)> { };
460
461 static inline size_t hash_thread_id(thread_id_t id)
462 {
463 static_assert(sizeof(thread_id_t) <= 8, "Expected a platform where thread IDs are at most 64-bit values");
464 return static_cast<size_t>(hash_32_or_64<sizeof(thread_id_converter<thread_id_t>::thread_id_hash_t)>::hash(
466 }
467
468 template<typename T>
469 static inline bool circular_less_than(T a, T b)
470 {
471#ifdef _MSC_VER
472#pragma warning(push)
473#pragma warning(disable: 4554)
474#endif
475 static_assert(std::is_integral<T>::value && !std::numeric_limits<T>::is_signed, "circular_less_than is intended to be used only with unsigned integer types");
476 return static_cast<T>(a - b) > static_cast<T>(static_cast<T>(1) << static_cast<T>(sizeof(T) * CHAR_BIT - 1));
477#ifdef _MSC_VER
478#pragma warning(pop)
479#endif
480 }
481
482 template<typename U>
483 static inline char* align_for(char* ptr)
484 {
485 const std::size_t alignment = std::alignment_of<U>::value;
486 return ptr + (alignment - (reinterpret_cast<std::uintptr_t>(ptr) % alignment)) % alignment;
487 }
488
489 template<typename T>
490 static inline T ceil_to_pow_2(T x)
491 {
492 static_assert(std::is_integral<T>::value && !std::numeric_limits<T>::is_signed, "ceil_to_pow_2 is intended to be used only with unsigned integer types");
493
494 // Adapted from http://graphics.stanford.edu/~seander/bithacks.html#RoundUpPowerOf2
495 --x;
496 x |= x >> 1;
497 x |= x >> 2;
498 x |= x >> 4;
499 for (std::size_t i = 1; i < sizeof(T); i <<= 1) {
500 x |= x >> (i << 3);
501 }
502 ++x;
503 return x;
504 }
505
506 template<typename T>
507 static inline void swap_relaxed(std::atomic<T>& left, std::atomic<T>& right)
508 {
509 T temp = std::move(left.load(std::memory_order_relaxed));
510 left.store(std::move(right.load(std::memory_order_relaxed)), std::memory_order_relaxed);
511 right.store(std::move(temp), std::memory_order_relaxed);
512 }
513
514 template<typename T>
515 static inline T const& nomove(T const& x)
516 {
517 return x;
518 }
519
520 template<bool Enable>
522 {
523 template<typename T>
524 static inline T const& eval(T const& x)
525 {
526 return x;
527 }
528 };
529
530 template<>
531 struct nomove_if<false>
532 {
533 template<typename U>
534 static inline auto eval(U&& x)
535 -> decltype(std::forward<U>(x))
536 {
537 return std::forward<U>(x);
538 }
539 };
540
541 template<typename It>
542 static inline auto deref_noexcept(It& it) MOODYCAMEL_NOEXCEPT -> decltype(*it)
543 {
544 return *it;
545 }
546
547#if defined(__clang__) || !defined(__GNUC__) || __GNUC__ > 4 || (__GNUC__ == 4 && __GNUC_MINOR__ >= 8)
548 template<typename T> struct is_trivially_destructible : std::is_trivially_destructible<T> { };
549#else
550 template<typename T> struct is_trivially_destructible : std::has_trivial_destructor<T> { };
551#endif
552
553#ifdef MOODYCAMEL_CPP11_THREAD_LOCAL_SUPPORTED
554#ifdef MCDBGQ_USE_RELACY
557#else
558 struct ThreadExitListener
559 {
560 typedef void (*callback_t)(void*);
561 callback_t callback;
562 void* userData;
563
564 ThreadExitListener* next; // reserved for use by the ThreadExitNotifier
565 };
566
567
568 class ThreadExitNotifier
569 {
570 public:
571 static void subscribe(ThreadExitListener* listener)
572 {
573 auto& tlsInst = instance();
574 listener->next = tlsInst.tail;
575 tlsInst.tail = listener;
576 }
577
578 static void unsubscribe(ThreadExitListener* listener)
579 {
580 auto& tlsInst = instance();
581 ThreadExitListener** prev = &tlsInst.tail;
582 for (auto ptr = tlsInst.tail; ptr != nullptr; ptr = ptr->next) {
583 if (ptr == listener) {
584 *prev = ptr->next;
585 break;
586 }
587 prev = &ptr->next;
588 }
589 }
590
591 private:
592 ThreadExitNotifier() : tail(nullptr) { }
593 ThreadExitNotifier(ThreadExitNotifier const&) MOODYCAMEL_DELETE_FUNCTION;
594 ThreadExitNotifier& operator=(ThreadExitNotifier const&) MOODYCAMEL_DELETE_FUNCTION;
595
596 ~ThreadExitNotifier()
597 {
598 // This thread is about to exit, let everyone know!
599 assert(this == &instance() && "If this assert fails, you likely have a buggy compiler! Change the preprocessor conditions such that MOODYCAMEL_CPP11_THREAD_LOCAL_SUPPORTED is no longer defined.");
600 for (auto ptr = tail; ptr != nullptr; ptr = ptr->next) {
601 ptr->callback(ptr->userData);
602 }
603 }
604
605 // Thread-local
606 static inline ThreadExitNotifier& instance()
607 {
608 static thread_local ThreadExitNotifier notifier;
609 return notifier;
610 }
611
612 private:
613 ThreadExitListener* tail;
614 };
615#endif
616#endif
617
618 template<typename T> struct static_is_lock_free_num { enum { value = 0 }; };
619 template<> struct static_is_lock_free_num<signed char> { enum { value = ATOMIC_CHAR_LOCK_FREE }; };
620 template<> struct static_is_lock_free_num<short> { enum { value = ATOMIC_SHORT_LOCK_FREE }; };
621 template<> struct static_is_lock_free_num<int> { enum { value = ATOMIC_INT_LOCK_FREE }; };
622 template<> struct static_is_lock_free_num<long> { enum { value = ATOMIC_LONG_LOCK_FREE }; };
623 template<> struct static_is_lock_free_num<long long> { enum { value = ATOMIC_LLONG_LOCK_FREE }; };
624 template<typename T> struct static_is_lock_free : static_is_lock_free_num<typename std::make_signed<T>::type> { };
625 template<> struct static_is_lock_free<bool> { enum { value = ATOMIC_BOOL_LOCK_FREE }; };
626 template<typename U> struct static_is_lock_free<U*> { enum { value = ATOMIC_POINTER_LOCK_FREE }; };
627}
628
629
631{
632 template<typename T, typename Traits>
634
635 template<typename T, typename Traits>
637
638 ProducerToken(ProducerToken&& other) MOODYCAMEL_NOEXCEPT
639 : producer(other.producer)
640 {
641 other.producer = nullptr;
642 if (producer != nullptr) {
643 producer->token = this;
644 }
645 }
646
647 inline ProducerToken& operator=(ProducerToken&& other) MOODYCAMEL_NOEXCEPT
648 {
649 swap(other);
650 return *this;
651 }
652
653 void swap(ProducerToken& other) MOODYCAMEL_NOEXCEPT
654 {
655 std::swap(producer, other.producer);
656 if (producer != nullptr) {
657 producer->token = this;
658 }
659 if (other.producer != nullptr) {
660 other.producer->token = &other;
661 }
662 }
663
664 // A token is always valid unless:
665 // 1) Memory allocation failed during construction
666 // 2) It was moved via the move constructor
667 // (Note: assignment does a swap, leaving both potentially valid)
668 // 3) The associated queue was destroyed
669 // Note that if valid() returns true, that only indicates
670 // that the token is valid for use with a specific queue,
671 // but not which one; that's up to the user to track.
672 inline bool valid() const { return producer != nullptr; }
673
675 {
676 if (producer != nullptr) {
677 producer->token = nullptr;
678 producer->inactive.store(true, std::memory_order_release);
679 }
680 }
681
682 // Disable copying and assignment
683 ProducerToken(ProducerToken const&) MOODYCAMEL_DELETE_FUNCTION;
684 ProducerToken& operator=(ProducerToken const&) MOODYCAMEL_DELETE_FUNCTION;
685
686private:
687 template<typename T, typename Traits> friend class ConcurrentQueue;
688 friend class ConcurrentQueueTests;
689
690protected:
692};
693
694
696{
697 template<typename T, typename Traits>
699
700 template<typename T, typename Traits>
702
703 ConsumerToken(ConsumerToken&& other) MOODYCAMEL_NOEXCEPT
704 : initialOffset(other.initialOffset), lastKnownGlobalOffset(other.lastKnownGlobalOffset), itemsConsumedFromCurrent(other.itemsConsumedFromCurrent), currentProducer(other.currentProducer), desiredProducer(other.desiredProducer)
705 {
706 }
707
708 inline ConsumerToken& operator=(ConsumerToken&& other) MOODYCAMEL_NOEXCEPT
709 {
710 swap(other);
711 return *this;
712 }
713
714 void swap(ConsumerToken& other) MOODYCAMEL_NOEXCEPT
715 {
716 std::swap(initialOffset, other.initialOffset);
717 std::swap(lastKnownGlobalOffset, other.lastKnownGlobalOffset);
718 std::swap(itemsConsumedFromCurrent, other.itemsConsumedFromCurrent);
719 std::swap(currentProducer, other.currentProducer);
720 std::swap(desiredProducer, other.desiredProducer);
721 }
722
723 // Disable copying and assignment
724 ConsumerToken(ConsumerToken const&) MOODYCAMEL_DELETE_FUNCTION;
725 ConsumerToken& operator=(ConsumerToken const&) MOODYCAMEL_DELETE_FUNCTION;
726
727private:
728 template<typename T, typename Traits> friend class ConcurrentQueue;
729 friend class ConcurrentQueueTests;
730
731private: // but shared with ConcurrentQueue
732 std::uint32_t initialOffset;
733 std::uint32_t lastKnownGlobalOffset;
734 std::uint32_t itemsConsumedFromCurrent;
737};
738
739// Need to forward-declare this swap because it's in a namespace.
740// See http://stackoverflow.com/questions/4492062/why-does-a-c-friend-class-need-a-forward-declaration-only-in-other-namespaces
741template<typename T, typename Traits>
742inline void swap(typename ConcurrentQueue<T, Traits>::ImplicitProducerKVP& a, typename ConcurrentQueue<T, Traits>::ImplicitProducerKVP& b) MOODYCAMEL_NOEXCEPT;
743
744
745template<typename T, typename Traits = ConcurrentQueueDefaultTraits>
747{
748public:
749 using value_type = T;
752
753 typedef typename Traits::index_t index_t;
754 typedef typename Traits::size_t size_t;
755
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);
762#ifdef _MSC_VER
763#pragma warning(push)
764#pragma warning(disable: 4307) // + integral constant overflow (that's what the ternary expression is for!)
765#pragma warning(disable: 4309) // static_cast: Truncation of constant value
766#endif
767 static const size_t MAX_SUBQUEUE_SIZE = (details::const_numeric_max<size_t>::value - static_cast<size_t>(Traits::MAX_SUBQUEUE_SIZE) < BLOCK_SIZE) ? details::const_numeric_max<size_t>::value : ((static_cast<size_t>(Traits::MAX_SUBQUEUE_SIZE) + (BLOCK_SIZE - 1)) / BLOCK_SIZE * BLOCK_SIZE);
768#ifdef _MSC_VER
769#pragma warning(pop)
770#endif
771
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)");
781
782public:
783 // Creates a queue with at least `capacity` element slots; note that the
784 // actual number of elements that can be inserted without additional memory
785 // allocation depends on the number of producers and the block size (e.g. if
786 // the block size is equal to `capacity`, only a single block will be allocated
787 // up-front, which means only a single producer will be able to enqueue elements
788 // without an extra allocation -- blocks aren't shared between producers).
789 // This method is not thread safe -- it is up to the user to ensure that the
790 // queue is fully constructed before it starts being used by other threads (this
791 // includes making the memory effects of construction visible, possibly with a
792 // memory barrier).
793 explicit ConcurrentQueue(size_t capacity = 6 * BLOCK_SIZE)
794 : producerListTail(nullptr),
795 producerCount(0),
796 initialBlockPoolIndex(0),
797 nextExplicitConsumerId(0),
798 globalExplicitConsumerOffset(0)
799 {
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));
803
804#ifdef MOODYCAMEL_QUEUE_INTERNAL_DEBUG
805 // Track all the producers using a fully-resolved typed list for
806 // each kind; this makes it possible to debug them starting from
807 // the root queue object (otherwise wacky casts are needed that
808 // don't compile in the debugger's expression evaluator).
809 explicitProducers.store(nullptr, std::memory_order_relaxed);
810 implicitProducers.store(nullptr, std::memory_order_relaxed);
811#endif
812 }
813
814 // Computes the correct amount of pre-allocated blocks for you based
815 // on the minimum number of elements you want available at any given
816 // time, and the maximum concurrent number of each type of producer.
818 : producerListTail(nullptr),
819 producerCount(0),
820 initialBlockPoolIndex(0),
821 nextExplicitConsumerId(0),
822 globalExplicitConsumerOffset(0)
823 {
824 implicitProducerHashResizeInProgress.clear(std::memory_order_relaxed);
825 populate_initial_implicit_producer_hash();
826 size_t blocks = (((minCapacity + BLOCK_SIZE - 1) / BLOCK_SIZE) - 1) * (maxExplicitProducers + 1) + 2 * (maxExplicitProducers + maxImplicitProducers);
827 populate_initial_block_list(blocks);
828
829#ifdef MOODYCAMEL_QUEUE_INTERNAL_DEBUG
830 explicitProducers.store(nullptr, std::memory_order_relaxed);
831 implicitProducers.store(nullptr, std::memory_order_relaxed);
832#endif
833 }
834
835 // Note: The queue should not be accessed concurrently while it's
836 // being deleted. It's up to the user to synchronize this.
837 // This method is not thread safe.
839 {
840 // Destroy producers
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;
846 }
847 destroy(ptr);
848 ptr = next;
849 }
850
851 // Destroy implicit producer hash tables
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) { // The last hash is part of this object and was not allocated dynamically
857 for (size_t i = 0; i != hash->capacity; ++i) {
858 hash->entries[i].~ImplicitProducerKVP();
859 }
860 hash->~ImplicitProducerHash();
861 (Traits::free)(hash);
862 }
863 hash = prev;
864 }
865 }
866
867 // Destroy global free list
868 auto block = freeList.head_unsafe();
869 while (block != nullptr) {
870 auto next = block->freeListNext.load(std::memory_order_relaxed);
871 if (block->dynamicallyAllocated) {
872 destroy(block);
873 }
874 block = next;
875 }
876
877 // Destroy initial free list
878 destroy_array(initialBlockPool, initialBlockPoolSize);
879 }
880
881 // Disable copying and copy assignment
882 ConcurrentQueue(ConcurrentQueue const&) MOODYCAMEL_DELETE_FUNCTION;
883 ConcurrentQueue& operator=(ConcurrentQueue const&) MOODYCAMEL_DELETE_FUNCTION;
884
885 // Moving is supported, but note that it is *not* a thread-safe operation.
886 // Nobody can use the queue while it's being moved, and the memory effects
887 // of that move must be propagated to other threads before they can use it.
888 // Note: When a queue is moved, its tokens are still valid but can only be
889 // used with the destination queue (i.e. semantically they are moved along
890 // with the queue itself).
891 ConcurrentQueue(ConcurrentQueue&& other) MOODYCAMEL_NOEXCEPT
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))
900 {
901 // Move the other one into this, and leave the other one as an empty queue
902 implicitProducerHashResizeInProgress.clear(std::memory_order_relaxed);
903 populate_initial_implicit_producer_hash();
904 swap_implicit_producer_hashes(other);
905
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);
910
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);
916#endif
917
918 other.initialBlockPoolIndex.store(0, std::memory_order_relaxed);
919 other.initialBlockPoolSize = 0;
920 other.initialBlockPool = nullptr;
921
922 reown_producers();
923 }
924
925 inline ConcurrentQueue& operator=(ConcurrentQueue&& other) MOODYCAMEL_NOEXCEPT
926 {
927 return swap_internal(other);
928 }
929
930 // Swaps this queue's state with the other's. Not thread-safe.
931 // Swapping two queues does not invalidate their tokens, however
932 // the tokens that were created for one queue must be used with
933 // only the swapped queue (i.e. the tokens are tied to the
934 // queue's movable state, not the object itself).
935 inline void swap(ConcurrentQueue& other) MOODYCAMEL_NOEXCEPT
936 {
937 swap_internal(other);
938 }
939
940private:
941 ConcurrentQueue& swap_internal(ConcurrentQueue& other)
942 {
943 if (this == &other) {
944 return *this;
945 }
946
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);
955
956 swap_implicit_producer_hashes(other);
957
958 reown_producers();
959 other.reown_producers();
960
961#ifdef MOODYCAMEL_QUEUE_INTERNAL_DEBUG
962 details::swap_relaxed(explicitProducers, other.explicitProducers);
963 details::swap_relaxed(implicitProducers, other.implicitProducers);
964#endif
965
966 return *this;
967 }
968
969public:
970 // Enqueues a single item (by copying it).
971 // Allocates memory if required. Only fails if memory allocation fails (or implicit
972 // production is disabled because Traits::INITIAL_IMPLICIT_PRODUCER_HASH_SIZE is 0,
973 // or Traits::MAX_SUBQUEUE_SIZE has been defined and would be surpassed).
974 // Thread-safe.
975 inline bool enqueue(T const& item)
976 {
977 MOODYCAMEL_CONSTEXPR_IF (INITIAL_IMPLICIT_PRODUCER_HASH_SIZE == 0) return false;
978 else return inner_enqueue<CanAlloc>(item);
979 }
980
981 // Enqueues a single item (by moving it, if possible).
982 // Allocates memory if required. Only fails if memory allocation fails (or implicit
983 // production is disabled because Traits::INITIAL_IMPLICIT_PRODUCER_HASH_SIZE is 0,
984 // or Traits::MAX_SUBQUEUE_SIZE has been defined and would be surpassed).
985 // Thread-safe.
986 inline bool enqueue(T&& item)
987 {
988 MOODYCAMEL_CONSTEXPR_IF (INITIAL_IMPLICIT_PRODUCER_HASH_SIZE == 0) return false;
989 else return inner_enqueue<CanAlloc>(std::move(item));
990 }
991
992 // Enqueues a single item (by copying it) using an explicit producer token.
993 // Allocates memory if required. Only fails if memory allocation fails (or
994 // Traits::MAX_SUBQUEUE_SIZE has been defined and would be surpassed).
995 // Thread-safe.
996 inline bool enqueue(producer_token_t const& token, T const& item)
997 {
998 return inner_enqueue<CanAlloc>(token, item);
999 }
1000
1001 // Enqueues a single item (by moving it, if possible) using an explicit producer token.
1002 // Allocates memory if required. Only fails if memory allocation fails (or
1003 // Traits::MAX_SUBQUEUE_SIZE has been defined and would be surpassed).
1004 // Thread-safe.
1005 inline bool enqueue(producer_token_t const& token, T&& item)
1006 {
1007 return inner_enqueue<CanAlloc>(token, std::move(item));
1008 }
1009
1010 // Enqueues several items.
1011 // Allocates memory if required. Only fails if memory allocation fails (or
1012 // implicit production is disabled because Traits::INITIAL_IMPLICIT_PRODUCER_HASH_SIZE
1013 // is 0, or Traits::MAX_SUBQUEUE_SIZE has been defined and would be surpassed).
1014 // Note: Use std::make_move_iterator if the elements should be moved instead of copied.
1015 // Thread-safe.
1016 template<typename It>
1017 bool enqueue_bulk(It itemFirst, size_t count)
1018 {
1019 MOODYCAMEL_CONSTEXPR_IF (INITIAL_IMPLICIT_PRODUCER_HASH_SIZE == 0) return false;
1020 else return inner_enqueue_bulk<CanAlloc>(itemFirst, count);
1021 }
1022
1023 // Enqueues several items using an explicit producer token.
1024 // Allocates memory if required. Only fails if memory allocation fails
1025 // (or Traits::MAX_SUBQUEUE_SIZE has been defined and would be surpassed).
1026 // Note: Use std::make_move_iterator if the elements should be moved
1027 // instead of copied.
1028 // Thread-safe.
1029 template<typename It>
1030 bool enqueue_bulk(producer_token_t const& token, It itemFirst, size_t count)
1031 {
1032 return inner_enqueue_bulk<CanAlloc>(token, itemFirst, count);
1033 }
1034
1035 // Enqueues a single item (by copying it).
1036 // Does not allocate memory. Fails if not enough room to enqueue (or implicit
1037 // production is disabled because Traits::INITIAL_IMPLICIT_PRODUCER_HASH_SIZE
1038 // is 0).
1039 // Thread-safe.
1040 inline bool try_enqueue(T const& item)
1041 {
1042 MOODYCAMEL_CONSTEXPR_IF (INITIAL_IMPLICIT_PRODUCER_HASH_SIZE == 0) return false;
1043 else return inner_enqueue<CannotAlloc>(item);
1044 }
1045
1046 // Enqueues a single item (by moving it, if possible).
1047 // Does not allocate memory (except for one-time implicit producer).
1048 // Fails if not enough room to enqueue (or implicit production is
1049 // disabled because Traits::INITIAL_IMPLICIT_PRODUCER_HASH_SIZE is 0).
1050 // Thread-safe.
1051 inline bool try_enqueue(T&& item)
1052 {
1053 MOODYCAMEL_CONSTEXPR_IF (INITIAL_IMPLICIT_PRODUCER_HASH_SIZE == 0) return false;
1054 else return inner_enqueue<CannotAlloc>(std::move(item));
1055 }
1056
1057 // Enqueues a single item (by copying it) using an explicit producer token.
1058 // Does not allocate memory. Fails if not enough room to enqueue.
1059 // Thread-safe.
1060 inline bool try_enqueue(producer_token_t const& token, T const& item)
1061 {
1062 return inner_enqueue<CannotAlloc>(token, item);
1063 }
1064
1065 // Enqueues a single item (by moving it, if possible) using an explicit producer token.
1066 // Does not allocate memory. Fails if not enough room to enqueue.
1067 // Thread-safe.
1068 inline bool try_enqueue(producer_token_t const& token, T&& item)
1069 {
1070 return inner_enqueue<CannotAlloc>(token, std::move(item));
1071 }
1072
1073 // Enqueues several items.
1074 // Does not allocate memory (except for one-time implicit producer).
1075 // Fails if not enough room to enqueue (or implicit production is
1076 // disabled because Traits::INITIAL_IMPLICIT_PRODUCER_HASH_SIZE is 0).
1077 // Note: Use std::make_move_iterator if the elements should be moved
1078 // instead of copied.
1079 // Thread-safe.
1080 template<typename It>
1081 bool try_enqueue_bulk(It itemFirst, size_t count)
1082 {
1083 MOODYCAMEL_CONSTEXPR_IF (INITIAL_IMPLICIT_PRODUCER_HASH_SIZE == 0) return false;
1084 else return inner_enqueue_bulk<CannotAlloc>(itemFirst, count);
1085 }
1086
1087 // Enqueues several items using an explicit producer token.
1088 // Does not allocate memory. Fails if not enough room to enqueue.
1089 // Note: Use std::make_move_iterator if the elements should be moved
1090 // instead of copied.
1091 // Thread-safe.
1092 template<typename It>
1093 bool try_enqueue_bulk(producer_token_t const& token, It itemFirst, size_t count)
1094 {
1095 return inner_enqueue_bulk<CannotAlloc>(token, itemFirst, count);
1096 }
1097
1098
1099
1100 // Attempts to dequeue from the queue.
1101 // Returns false if all producer streams appeared empty at the time they
1102 // were checked (so, the queue is likely but not guaranteed to be empty).
1103 // Never allocates. Thread-safe.
1104 template<typename U>
1105 bool try_dequeue(U& item)
1106 {
1107 // Instead of simply trying each producer in turn (which could cause needless contention on the first
1108 // producer), we score them heuristically.
1109 size_t nonEmptyCount = 0;
1110 ProducerBase* best = nullptr;
1111 size_t bestSize = 0;
1112 for (auto ptr = producerListTail.load(std::memory_order_acquire); nonEmptyCount < 3 && ptr != nullptr; ptr = ptr->next_prod()) {
1113 auto size = ptr->size_approx();
1114 if (size > 0) {
1115 if (size > bestSize) {
1116 bestSize = size;
1117 best = ptr;
1118 }
1119 ++nonEmptyCount;
1120 }
1121 }
1122
1123 // If there was at least one non-empty queue but it appears empty at the time
1124 // we try to dequeue from it, we need to make sure every queue's been tried
1125 if (nonEmptyCount > 0) {
1126 if ((details::likely)(best->dequeue(item))) {
1127 return true;
1128 }
1129 for (auto ptr = producerListTail.load(std::memory_order_acquire); ptr != nullptr; ptr = ptr->next_prod()) {
1130 if (ptr != best && ptr->dequeue(item)) {
1131 return true;
1132 }
1133 }
1134 }
1135 return false;
1136 }
1137
1138 // Attempts to dequeue from the queue.
1139 // Returns false if all producer streams appeared empty at the time they
1140 // were checked (so, the queue is likely but not guaranteed to be empty).
1141 // This differs from the try_dequeue(item) method in that this one does
1142 // not attempt to reduce contention by interleaving the order that producer
1143 // streams are dequeued from. So, using this method can reduce overall throughput
1144 // under contention, but will give more predictable results in single-threaded
1145 // consumer scenarios. This is mostly only useful for internal unit tests.
1146 // Never allocates. Thread-safe.
1147 template<typename U>
1148 bool try_dequeue_non_interleaved(U& item)
1149 {
1150 for (auto ptr = producerListTail.load(std::memory_order_acquire); ptr != nullptr; ptr = ptr->next_prod()) {
1151 if (ptr->dequeue(item)) {
1152 return true;
1153 }
1154 }
1155 return false;
1156 }
1157
1158 // Attempts to dequeue from the queue using an explicit consumer token.
1159 // Returns false if all producer streams appeared empty at the time they
1160 // were checked (so, the queue is likely but not guaranteed to be empty).
1161 // Never allocates. Thread-safe.
1162 template<typename U>
1163 bool try_dequeue(consumer_token_t& token, U& item)
1164 {
1165 // The idea is roughly as follows:
1166 // Every 256 items from one producer, make everyone rotate (increase the global offset) -> this means the highest efficiency consumer dictates the rotation speed of everyone else, more or less
1167 // If you see that the global offset has changed, you must reset your consumption counter and move to your designated place
1168 // If there's no items where you're supposed to be, keep moving until you find a producer with some items
1169 // If the global offset has not changed but you've run out of items to consume, move over from your current position until you find an producer with something in it
1170
1171 if (token.desiredProducer == nullptr || token.lastKnownGlobalOffset != globalExplicitConsumerOffset.load(std::memory_order_relaxed)) {
1172 if (!update_current_producer_after_rotation(token)) {
1173 return false;
1174 }
1175 }
1176
1177 // If there was at least one non-empty queue but it appears empty at the time
1178 // we try to dequeue from it, we need to make sure every queue's been tried
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);
1182 }
1183 return true;
1184 }
1185
1186 auto tail = producerListTail.load(std::memory_order_acquire);
1187 auto ptr = static_cast<ProducerBase*>(token.currentProducer)->next_prod();
1188 if (ptr == nullptr) {
1189 ptr = tail;
1190 }
1191 while (ptr != static_cast<ProducerBase*>(token.currentProducer)) {
1192 if (ptr->dequeue(item)) {
1193 token.currentProducer = ptr;
1194 token.itemsConsumedFromCurrent = 1;
1195 return true;
1196 }
1197 ptr = ptr->next_prod();
1198 if (ptr == nullptr) {
1199 ptr = tail;
1200 }
1201 }
1202 return false;
1203 }
1204
1205 // Attempts to dequeue several elements from the queue.
1206 // Returns the number of items actually dequeued.
1207 // Returns 0 if all producer streams appeared empty at the time they
1208 // were checked (so, the queue is likely but not guaranteed to be empty).
1209 // Never allocates. Thread-safe.
1210 template<typename It>
1211 size_t try_dequeue_bulk(It itemFirst, size_t max)
1212 {
1213 size_t count = 0;
1214 for (auto ptr = producerListTail.load(std::memory_order_acquire); ptr != nullptr; ptr = ptr->next_prod()) {
1215 count += ptr->dequeue_bulk(itemFirst, max - count);
1216 if (count == max) {
1217 break;
1218 }
1219 }
1220 return count;
1221 }
1222
1223 // Attempts to dequeue several elements from the queue using an explicit consumer token.
1224 // Returns the number of items actually dequeued.
1225 // Returns 0 if all producer streams appeared empty at the time they
1226 // were checked (so, the queue is likely but not guaranteed to be empty).
1227 // Never allocates. Thread-safe.
1228 template<typename It>
1229 size_t try_dequeue_bulk(consumer_token_t& token, It itemFirst, size_t max)
1230 {
1231 if (token.desiredProducer == nullptr || token.lastKnownGlobalOffset != globalExplicitConsumerOffset.load(std::memory_order_relaxed)) {
1232 if (!update_current_producer_after_rotation(token)) {
1233 return 0;
1234 }
1235 }
1236
1237 size_t count = static_cast<ProducerBase*>(token.currentProducer)->dequeue_bulk(itemFirst, max);
1238 if (count == 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);
1241 }
1242 return max;
1243 }
1244 token.itemsConsumedFromCurrent += static_cast<std::uint32_t>(count);
1245 max -= count;
1246
1247 auto tail = producerListTail.load(std::memory_order_acquire);
1248 auto ptr = static_cast<ProducerBase*>(token.currentProducer)->next_prod();
1249 if (ptr == nullptr) {
1250 ptr = tail;
1251 }
1252 while (ptr != static_cast<ProducerBase*>(token.currentProducer)) {
1253 auto dequeued = ptr->dequeue_bulk(itemFirst, max);
1254 count += dequeued;
1255 if (dequeued != 0) {
1256 token.currentProducer = ptr;
1257 token.itemsConsumedFromCurrent = static_cast<std::uint32_t>(dequeued);
1258 }
1259 if (dequeued == max) {
1260 break;
1261 }
1262 max -= dequeued;
1263 ptr = ptr->next_prod();
1264 if (ptr == nullptr) {
1265 ptr = tail;
1266 }
1267 }
1268 return count;
1269 }
1270
1271
1272
1273 // Attempts to dequeue from a specific producer's inner queue.
1274 // If you happen to know which producer you want to dequeue from, this
1275 // is significantly faster than using the general-case try_dequeue methods.
1276 // Returns false if the producer's queue appeared empty at the time it
1277 // was checked (so, the queue is likely but not guaranteed to be empty).
1278 // Never allocates. Thread-safe.
1279 template<typename U>
1280 inline bool try_dequeue_from_producer(producer_token_t const& producer, U& item)
1281 {
1282 return static_cast<ExplicitProducer*>(producer.producer)->dequeue(item);
1283 }
1284
1285 // Attempts to dequeue several elements from a specific producer's inner queue.
1286 // Returns the number of items actually dequeued.
1287 // If you happen to know which producer you want to dequeue from, this
1288 // is significantly faster than using the general-case try_dequeue methods.
1289 // Returns 0 if the producer's queue appeared empty at the time it
1290 // was checked (so, the queue is likely but not guaranteed to be empty).
1291 // Never allocates. Thread-safe.
1292 template<typename It>
1293 inline size_t try_dequeue_bulk_from_producer(producer_token_t const& producer, It itemFirst, size_t max)
1294 {
1295 return static_cast<ExplicitProducer*>(producer.producer)->dequeue_bulk(itemFirst, max);
1296 }
1297
1298
1299 // Returns an estimate of the total number of elements currently in the queue. This
1300 // estimate is only accurate if the queue has completely stabilized before it is called
1301 // (i.e. all enqueue and dequeue operations have completed and their memory effects are
1302 // visible on the calling thread, and no further operations start while this method is
1303 // being called).
1304 // Thread-safe.
1305 size_t size_approx() const
1306 {
1307 size_t size = 0;
1308 for (auto ptr = producerListTail.load(std::memory_order_acquire); ptr != nullptr; ptr = ptr->next_prod()) {
1309 size += ptr->size_approx();
1310 }
1311 return size;
1312 }
1313
1314
1315 // Returns true if the underlying atomic variables used by
1316 // the queue are lock-free (they should be on most platforms).
1317 // Thread-safe.
1318 static bool is_lock_free()
1319 {
1320 return
1327 }
1328
1329
1330private:
1331 friend struct ProducerToken;
1332 friend struct ConsumerToken;
1333 struct ExplicitProducer;
1334 friend struct ExplicitProducer;
1335 struct ImplicitProducer;
1336 friend struct ImplicitProducer;
1337 friend class ConcurrentQueueTests;
1338
1339 enum AllocationMode { CanAlloc, CannotAlloc };
1340
1341
1343 // Queue methods
1345
1346 template<AllocationMode canAlloc, typename U>
1347 inline bool inner_enqueue(producer_token_t const& token, U&& element)
1348 {
1349 return static_cast<ExplicitProducer*>(token.producer)->ConcurrentQueue::ExplicitProducer::template enqueue<canAlloc>(std::forward<U>(element));
1350 }
1351
1352 template<AllocationMode canAlloc, typename U>
1353 inline bool inner_enqueue(U&& element)
1354 {
1355 auto producer = get_or_add_implicit_producer();
1356 return producer == nullptr ? false : producer->ConcurrentQueue::ImplicitProducer::template enqueue<canAlloc>(std::forward<U>(element));
1357 }
1358
1359 template<AllocationMode canAlloc, typename It>
1360 inline bool inner_enqueue_bulk(producer_token_t const& token, It itemFirst, size_t count)
1361 {
1362 return static_cast<ExplicitProducer*>(token.producer)->ConcurrentQueue::ExplicitProducer::template enqueue_bulk<canAlloc>(itemFirst, count);
1363 }
1364
1365 template<AllocationMode canAlloc, typename It>
1366 inline bool inner_enqueue_bulk(It itemFirst, size_t count)
1367 {
1368 auto producer = get_or_add_implicit_producer();
1369 return producer == nullptr ? false : producer->ConcurrentQueue::ImplicitProducer::template enqueue_bulk<canAlloc>(itemFirst, count);
1370 }
1371
1372 inline bool update_current_producer_after_rotation(consumer_token_t& token)
1373 {
1374 // Ah, there's been a rotation, figure out where we should be!
1375 auto tail = producerListTail.load(std::memory_order_acquire);
1376 if (token.desiredProducer == nullptr && tail == nullptr) {
1377 return false;
1378 }
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)) {
1382 // Aha, first time we're dequeueing anything.
1383 // Figure out our local position
1384 // Note: offset is from start, not end, but we're traversing from end -- subtract from count first
1385 std::uint32_t offset = prodCount - 1 - (token.initialOffset % prodCount);
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;
1391 }
1392 }
1393 }
1394
1395 std::uint32_t delta = globalOffset - token.lastKnownGlobalOffset;
1396 if (delta >= prodCount) {
1397 delta = delta % prodCount;
1398 }
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;
1403 }
1404 }
1405
1406 token.lastKnownGlobalOffset = globalOffset;
1407 token.currentProducer = token.desiredProducer;
1408 token.itemsConsumedFromCurrent = 0;
1409 return true;
1410 }
1411
1412
1414 // Free list
1416
1417 template <typename N>
1418 struct FreeListNode
1419 {
1420 FreeListNode() : freeListRefs(0), freeListNext(nullptr) { }
1421
1422 std::atomic<std::uint32_t> freeListRefs;
1423 std::atomic<N*> freeListNext;
1424 };
1425
1426 // A simple CAS-based lock-free free list. Not the fastest thing in the world under heavy contention, but
1427 // simple and correct (assuming nodes are never freed until after the free list is destroyed), and fairly
1428 // speedy under low contention.
1429 template<typename N> // N must inherit FreeListNode or have the same fields (and initialization of them)
1430 struct FreeList
1431 {
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); }
1435
1436 FreeList(FreeList const&) MOODYCAMEL_DELETE_FUNCTION;
1437 FreeList& operator=(FreeList const&) MOODYCAMEL_DELETE_FUNCTION;
1438
1439 inline void add(N* node)
1440 {
1441#ifdef MCDBGQ_NOLOCKFREE_FREELIST
1442 debug::DebugLock lock(mutex);
1443#endif
1444 // We know that the should-be-on-freelist bit is 0 at this point, so it's safe to
1445 // set it using a fetch_add
1446 if (node->freeListRefs.fetch_add(SHOULD_BE_ON_FREELIST, std::memory_order_acq_rel) == 0) {
1447 // Oh look! We were the last ones referencing this node, and we know
1448 // we want to add it to the free list, so let's do it!
1449 add_knowing_refcount_is_zero(node);
1450 }
1451 }
1452
1453 inline N* try_get()
1454 {
1455#ifdef MCDBGQ_NOLOCKFREE_FREELIST
1456 debug::DebugLock lock(mutex);
1457#endif
1458 auto head = freeListHead.load(std::memory_order_acquire);
1459 while (head != nullptr) {
1460 auto prevHead = head;
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);
1464 continue;
1465 }
1466
1467 // Good, reference count has been incremented (it wasn't at zero), which means we can read the
1468 // next and not worry about it changing between now and the time we do the CAS
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)) {
1471 // Yay, got the node. This means it was on the list, which means shouldBeOnFreeList must be false no
1472 // matter the refcount (because nobody else knows it's been taken off yet, it can't have been put back on).
1473 assert((head->freeListRefs.load(std::memory_order_relaxed) & SHOULD_BE_ON_FREELIST) == 0);
1474
1475 // Decrease refcount twice, once for our ref, and once for the list's ref
1476 head->freeListRefs.fetch_sub(2, std::memory_order_release);
1477 return head;
1478 }
1479
1480 // OK, the head must have changed on us, but we still need to decrease the refcount we increased.
1481 // Note that we don't need to release any memory effects, but we do need to ensure that the reference
1482 // count decrement happens-after the CAS on the head.
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);
1486 }
1487 }
1488
1489 return nullptr;
1490 }
1491
1492 // Useful for traversing the list when there's no contention (e.g. to destroy remaining nodes)
1493 N* head_unsafe() const { return freeListHead.load(std::memory_order_relaxed); }
1494
1495 private:
1496 inline void add_knowing_refcount_is_zero(N* node)
1497 {
1498 // Since the refcount is zero, and nobody can increase it once it's zero (except us, and we run
1499 // only one copy of this method per node at a time, i.e. the single thread case), then we know
1500 // we can safely change the next pointer of the node; however, once the refcount is back above
1501 // zero, then other threads could increase it (happens under heavy contention, when the refcount
1502 // goes to zero in between a load and a refcount increment of a node in try_get, then back up to
1503 // something non-zero, then the refcount increment is done by the other thread) -- so, if the CAS
1504 // to add the node to the actual list fails, decrease the refcount and leave the add operation to
1505 // the next thread who puts the refcount back at zero (which could be us, hence the loop).
1506 auto head = freeListHead.load(std::memory_order_relaxed);
1507 while (true) {
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)) {
1511 // Hmm, the add failed, but we can only try again when the refcount goes back to zero
1512 if (node->freeListRefs.fetch_add(SHOULD_BE_ON_FREELIST - 1, std::memory_order_release) == 1) {
1513 continue;
1514 }
1515 }
1516 return;
1517 }
1518 }
1519
1520 private:
1521 // Implemented like a stack, but where node order doesn't matter (nodes are inserted out of order under contention)
1522 std::atomic<N*> freeListHead;
1523
1524 static const std::uint32_t REFS_MASK = 0x7FFFFFFF;
1525 static const std::uint32_t SHOULD_BE_ON_FREELIST = 0x80000000;
1526
1527#ifdef MCDBGQ_NOLOCKFREE_FREELIST
1528 debug::DebugMutex mutex;
1529#endif
1530 };
1531
1532
1534 // Block
1536
1537 enum InnerQueueContext { implicit_context = 0, explicit_context = 1 };
1538
1539 struct Block
1540 {
1541 Block()
1542 : next(nullptr), elementsCompletelyDequeued(0), freeListRefs(0), freeListNext(nullptr), shouldBeOnFreeList(false), dynamicallyAllocated(true)
1543 {
1544#ifdef MCDBGQ_TRACKMEM
1545 owner = nullptr;
1546#endif
1547 }
1548
1549 template<InnerQueueContext context>
1550 inline bool is_empty() const
1551 {
1552 MOODYCAMEL_CONSTEXPR_IF (context == explicit_context && BLOCK_SIZE <= EXPLICIT_BLOCK_EMPTY_COUNTER_THRESHOLD) {
1553 // Check flags
1554 for (size_t i = 0; i < BLOCK_SIZE; ++i) {
1555 if (!emptyFlags[i].load(std::memory_order_relaxed)) {
1556 return false;
1557 }
1558 }
1559
1560 // Aha, empty; make sure we have all other memory effects that happened before the empty flags were set
1561 std::atomic_thread_fence(std::memory_order_acquire);
1562 return true;
1563 }
1564 else {
1565 // Check counter
1566 if (elementsCompletelyDequeued.load(std::memory_order_relaxed) == BLOCK_SIZE) {
1567 std::atomic_thread_fence(std::memory_order_acquire);
1568 return true;
1569 }
1570 assert(elementsCompletelyDequeued.load(std::memory_order_relaxed) <= BLOCK_SIZE);
1571 return false;
1572 }
1573 }
1574
1575 // Returns true if the block is now empty (does not apply in explicit context)
1576 template<InnerQueueContext context>
1577 inline bool set_empty(MOODYCAMEL_MAYBE_UNUSED index_t i)
1578 {
1579 MOODYCAMEL_CONSTEXPR_IF (context == explicit_context && BLOCK_SIZE <= EXPLICIT_BLOCK_EMPTY_COUNTER_THRESHOLD) {
1580 // Set flag
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);
1583 return false;
1584 }
1585 else {
1586 // Increment counter
1587 auto prevVal = elementsCompletelyDequeued.fetch_add(1, std::memory_order_release);
1588 assert(prevVal < BLOCK_SIZE);
1589 return prevVal == BLOCK_SIZE - 1;
1590 }
1591 }
1592
1593 // Sets multiple contiguous item statuses to 'empty' (assumes no wrapping and count > 0).
1594 // Returns true if the block is now empty (does not apply in explicit context).
1595 template<InnerQueueContext context>
1596 inline bool set_many_empty(MOODYCAMEL_MAYBE_UNUSED index_t i, size_t count)
1597 {
1598 MOODYCAMEL_CONSTEXPR_IF (context == explicit_context && BLOCK_SIZE <= EXPLICIT_BLOCK_EMPTY_COUNTER_THRESHOLD) {
1599 // Set flags
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);
1605 }
1606 return false;
1607 }
1608 else {
1609 // Increment counter
1610 auto prevVal = elementsCompletelyDequeued.fetch_add(count, std::memory_order_release);
1611 assert(prevVal + count <= BLOCK_SIZE);
1612 return prevVal + count == BLOCK_SIZE;
1613 }
1614 }
1615
1616 template<InnerQueueContext context>
1617 inline void set_all_empty()
1618 {
1619 MOODYCAMEL_CONSTEXPR_IF (context == explicit_context && BLOCK_SIZE <= EXPLICIT_BLOCK_EMPTY_COUNTER_THRESHOLD) {
1620 // Set all flags
1621 for (size_t i = 0; i != BLOCK_SIZE; ++i) {
1622 emptyFlags[i].store(true, std::memory_order_relaxed);
1623 }
1624 }
1625 else {
1626 // Reset counter
1627 elementsCompletelyDequeued.store(BLOCK_SIZE, std::memory_order_relaxed);
1628 }
1629 }
1630
1631 template<InnerQueueContext context>
1632 inline void reset_empty()
1633 {
1634 MOODYCAMEL_CONSTEXPR_IF (context == explicit_context && BLOCK_SIZE <= EXPLICIT_BLOCK_EMPTY_COUNTER_THRESHOLD) {
1635 // Reset flags
1636 for (size_t i = 0; i != BLOCK_SIZE; ++i) {
1637 emptyFlags[i].store(false, std::memory_order_relaxed);
1638 }
1639 }
1640 else {
1641 // Reset counter
1642 elementsCompletelyDequeued.store(0, std::memory_order_relaxed);
1643 }
1644 }
1645
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)); }
1648
1649 private:
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;
1652 public:
1653 Block* next;
1654 std::atomic<size_t> elementsCompletelyDequeued;
1655 std::atomic<bool> emptyFlags[BLOCK_SIZE <= EXPLICIT_BLOCK_EMPTY_COUNTER_THRESHOLD ? BLOCK_SIZE : 1];
1656 public:
1657 std::atomic<std::uint32_t> freeListRefs;
1658 std::atomic<Block*> freeListNext;
1659 std::atomic<bool> shouldBeOnFreeList;
1660 bool dynamicallyAllocated; // Perhaps a better name for this would be 'isNotPartOfInitialBlockPool'
1661
1662#ifdef MCDBGQ_TRACKMEM
1663 void* owner;
1664#endif
1665 };
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");
1667
1668
1669#ifdef MCDBGQ_TRACKMEM
1670public:
1671 struct MemStats;
1672private:
1673#endif
1674
1676 // Producer base
1678
1679 struct ProducerBase : public details::ConcurrentQueueProducerTypelessBase
1680 {
1681 ProducerBase(ConcurrentQueue* parent_, bool isExplicit_) :
1682 tailIndex(0),
1683 headIndex(0),
1684 dequeueOptimisticCount(0),
1685 dequeueOvercommit(0),
1686 tailBlock(nullptr),
1687 isExplicit(isExplicit_),
1688 parent(parent_)
1689 {
1690 }
1691
1692 virtual ~ProducerBase() { }
1693
1694 template<typename U>
1695 inline bool dequeue(U& element)
1696 {
1697 if (isExplicit) {
1698 return static_cast<ExplicitProducer*>(this)->dequeue(element);
1699 }
1700 else {
1701 return static_cast<ImplicitProducer*>(this)->dequeue(element);
1702 }
1703 }
1704
1705 template<typename It>
1706 inline size_t dequeue_bulk(It& itemFirst, size_t max)
1707 {
1708 if (isExplicit) {
1709 return static_cast<ExplicitProducer*>(this)->dequeue_bulk(itemFirst, max);
1710 }
1711 else {
1712 return static_cast<ImplicitProducer*>(this)->dequeue_bulk(itemFirst, max);
1713 }
1714 }
1715
1716 inline ProducerBase* next_prod() const { return static_cast<ProducerBase*>(next); }
1717
1718 inline size_t size_approx() const
1719 {
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;
1723 }
1724
1725 inline index_t getTail() const { return tailIndex.load(std::memory_order_relaxed); }
1726 protected:
1727 std::atomic<index_t> tailIndex; // Where to enqueue to next
1728 std::atomic<index_t> headIndex; // Where to dequeue from next
1729
1730 std::atomic<index_t> dequeueOptimisticCount;
1731 std::atomic<index_t> dequeueOvercommit;
1732
1733 Block* tailBlock;
1734
1735 public:
1736 bool isExplicit;
1737 ConcurrentQueue* parent;
1738
1739 protected:
1740#ifdef MCDBGQ_TRACKMEM
1741 friend struct MemStats;
1742#endif
1743 };
1744
1745
1747 // Explicit queue
1749
1750 struct ExplicitProducer : public ProducerBase
1751 {
1752 explicit ExplicitProducer(ConcurrentQueue* parent_) :
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)
1760 {
1761 size_t poolBasedIndexSize = details::ceil_to_pow_2(parent_->initialBlockPoolSize) >> 1;
1762 if (poolBasedIndexSize > pr_blockIndexSize) {
1763 pr_blockIndexSize = poolBasedIndexSize;
1764 }
1765
1766 new_block_index(0); // This creates an index with double the number of current entries, i.e. EXPLICIT_INITIAL_INDEX_SIZE
1767 }
1768
1769 ~ExplicitProducer()
1770 {
1771 // Destruct any elements not yet dequeued.
1772 // Since we're in the destructor, we can assume all elements
1773 // are either completely dequeued or completely not (no halfways).
1774 if (this->tailBlock != nullptr) { // Note this means there must be a block index too
1775 // First find the block that's partially dequeued, if any
1776 Block* halfDequeuedBlock = nullptr;
1777 if ((this->headIndex.load(std::memory_order_relaxed) & static_cast<index_t>(BLOCK_SIZE - 1)) != 0) {
1778 // The head's not on a block boundary, meaning a block somewhere is partially dequeued
1779 // (or the head block is the tail block and was fully dequeued, but the head/tail are still not on a boundary)
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);
1783 }
1784 assert(details::circular_less_than<index_t>(pr_blockIndexEntries[i].base, this->headIndex.load(std::memory_order_relaxed)));
1785 halfDequeuedBlock = pr_blockIndexEntries[i].block;
1786 }
1787
1788 // Start at the head block (note the first line in the loop gives us the head from the tail on the first iteration)
1789 auto block = this->tailBlock;
1790 do {
1791 block = block->next;
1792 if (block->ConcurrentQueue::Block::template is_empty<explicit_context>()) {
1793 continue;
1794 }
1795
1796 size_t i = 0; // Offset into block
1797 if (block == halfDequeuedBlock) {
1798 i = static_cast<size_t>(this->headIndex.load(std::memory_order_relaxed) & static_cast<index_t>(BLOCK_SIZE - 1));
1799 }
1800
1801 // Walk through all the items in the block; if this is the tail block, we need to stop when we reach the tail index
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();
1805 }
1806 } while (block != this->tailBlock);
1807 }
1808
1809 // Destroy all blocks that we own
1810 if (this->tailBlock != nullptr) {
1811 auto block = this->tailBlock;
1812 do {
1813 auto nextBlock = block->next;
1814 if (block->dynamicallyAllocated) {
1815 destroy(block);
1816 }
1817 else {
1818 this->parent->add_block_to_free_list(block);
1819 }
1820 block = nextBlock;
1821 } while (block != this->tailBlock);
1822 }
1823
1824 // Destroy the block indices
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);
1830 header = prev;
1831 }
1832 }
1833
1834 template<AllocationMode allocMode, typename U>
1835 inline bool enqueue(U&& element)
1836 {
1837 index_t currentTailIndex = this->tailIndex.load(std::memory_order_relaxed);
1838 index_t newTailIndex = 1 + currentTailIndex;
1839 if ((currentTailIndex & static_cast<index_t>(BLOCK_SIZE - 1)) == 0) {
1840 // We reached the end of a block, start a new one
1841 auto startBlock = this->tailBlock;
1842 auto originalBlockIndexSlotsUsed = pr_blockIndexSlotsUsed;
1843 if (this->tailBlock != nullptr && this->tailBlock->next->ConcurrentQueue::Block::template is_empty<explicit_context>()) {
1844 // We can re-use the block ahead of us, it's empty!
1845 this->tailBlock = this->tailBlock->next;
1846 this->tailBlock->ConcurrentQueue::Block::template reset_empty<explicit_context>();
1847
1848 // We'll put the block on the block index (guaranteed to be room since we're conceptually removing the
1849 // last block from it first -- except instead of removing then adding, we can just overwrite).
1850 // Note that there must be a valid block index here, since even if allocation failed in the ctor,
1851 // it would have been re-attempted when adding the first block to the queue; since there is such
1852 // a block, a block index must have been successfully allocated.
1853 }
1854 else {
1855 // Whatever head value we see here is >= the last value we saw here (relatively),
1856 // and <= its current value. Since we have the most recent tail, the head must be
1857 // <= to it.
1858 auto head = this->headIndex.load(std::memory_order_relaxed);
1859 assert(!details::circular_less_than<index_t>(currentTailIndex, head));
1860 if (!details::circular_less_than<index_t>(head, currentTailIndex + BLOCK_SIZE)
1861 || (MAX_SUBQUEUE_SIZE != details::const_numeric_max<size_t>::value && (MAX_SUBQUEUE_SIZE == 0 || MAX_SUBQUEUE_SIZE - BLOCK_SIZE < currentTailIndex - head))) {
1862 // We can't enqueue in another block because there's not enough leeway -- the
1863 // tail could surpass the head by the time the block fills up! (Or we'll exceed
1864 // the size limit, if the second part of the condition was true.)
1865 return false;
1866 }
1867 // We're going to need a new block; check that the block index has room
1868 if (pr_blockIndexRaw == nullptr || pr_blockIndexSlotsUsed == pr_blockIndexSize) {
1869 // Hmm, the circular block index is already full -- we'll need
1870 // to allocate a new index. Note pr_blockIndexRaw can only be nullptr if
1871 // the initial allocation failed in the constructor.
1872
1873 MOODYCAMEL_CONSTEXPR_IF (allocMode == CannotAlloc) {
1874 return false;
1875 }
1876 else if (!new_block_index(pr_blockIndexSlotsUsed)) {
1877 return false;
1878 }
1879 }
1880
1881 // Insert a new block in the circular linked list
1882 auto newBlock = this->parent->ConcurrentQueue::template requisition_block<allocMode>();
1883 if (newBlock == nullptr) {
1884 return false;
1885 }
1886#ifdef MCDBGQ_TRACKMEM
1887 newBlock->owner = this;
1888#endif
1889 newBlock->ConcurrentQueue::Block::template reset_empty<explicit_context>();
1890 if (this->tailBlock == nullptr) {
1891 newBlock->next = newBlock;
1892 }
1893 else {
1894 newBlock->next = this->tailBlock->next;
1895 this->tailBlock->next = newBlock;
1896 }
1897 this->tailBlock = newBlock;
1898 ++pr_blockIndexSlotsUsed;
1899 }
1900
1901 MOODYCAMEL_CONSTEXPR_IF (!MOODYCAMEL_NOEXCEPT_CTOR(T, U, new (static_cast<T*>(nullptr)) T(std::forward<U>(element)))) {
1902 // The constructor may throw. We want the element not to appear in the queue in
1903 // that case (without corrupting the queue):
1904 MOODYCAMEL_TRY {
1905 new ((*this->tailBlock)[currentTailIndex]) T(std::forward<U>(element));
1906 }
1907 MOODYCAMEL_CATCH (...) {
1908 // Revert change to the current block, but leave the new block available
1909 // for next time
1910 pr_blockIndexSlotsUsed = originalBlockIndexSlotsUsed;
1911 this->tailBlock = startBlock == nullptr ? this->tailBlock : startBlock;
1912 MOODYCAMEL_RETHROW;
1913 }
1914 }
1915 else {
1916 (void)startBlock;
1918 }
1919
1920 // Add block to block index
1921 auto& entry = blockIndex.load(std::memory_order_relaxed)->entries[pr_blockIndexFront];
1922 entry.base = currentTailIndex;
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);
1926
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);
1929 return true;
1930 }
1931 }
1932
1933 // Enqueue
1934 new ((*this->tailBlock)[currentTailIndex]) T(std::forward<U>(element));
1935
1936 this->tailIndex.store(newTailIndex, std::memory_order_release);
1937 return true;
1938 }
1939
1940 template<typename U>
1941 bool dequeue(U& element)
1942 {
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)) {
1946 // Might be something to dequeue, let's give it a try
1947
1948 // Note that this if is purely for performance purposes in the common case when the queue is
1949 // empty and the values are eventually consistent -- we may enter here spuriously.
1950
1951 // Note that whatever the values of overcommit and tail are, they are not going to change (unless we
1952 // change them) and must be the same value at this point (inside the if) as when the if condition was
1953 // evaluated.
1954
1955 // We insert an acquire fence here to synchronize-with the release upon incrementing dequeueOvercommit below.
1956 // This ensures that whatever the value we got loaded into overcommit, the load of dequeueOptisticCount in
1957 // the fetch_add below will result in a value at least as recent as that (and therefore at least as large).
1958 // Note that I believe a compiler (signal) fence here would be sufficient due to the nature of fetch_add (all
1959 // read-modify-write operations are guaranteed to work on the latest value in the modification order), but
1960 // unfortunately that can't be shown to be correct using only the C++11 standard.
1961 // See http://stackoverflow.com/questions/18223161/what-are-the-c11-memory-ordering-guarantees-in-this-corner-case
1962 std::atomic_thread_fence(std::memory_order_acquire);
1963
1964 // Increment optimistic counter, then check if it went over the boundary
1965 auto myDequeueCount = this->dequeueOptimisticCount.fetch_add(1, std::memory_order_relaxed);
1966
1967 // Note that since dequeueOvercommit must be <= dequeueOptimisticCount (because dequeueOvercommit is only ever
1968 // incremented after dequeueOptimisticCount -- this is enforced in the `else` block below), and since we now
1969 // have a version of dequeueOptimisticCount that is at least as recent as overcommit (due to the release upon
1970 // incrementing dequeueOvercommit and the acquire above that synchronizes with it), overcommit <= myDequeueCount.
1971 // However, we can't assert this since both dequeueOptimisticCount and dequeueOvercommit may (independently)
1972 // overflow; in such a case, though, the logic still holds since the difference between the two is maintained.
1973
1974 // Note that we reload tail here in case it changed; it will be the same value as before or greater, since
1975 // this load is sequenced after (happens after) the earlier load above. This is supported by read-read
1976 // coherency (as defined in the standard), explained here: http://en.cppreference.com/w/cpp/atomic/memory_order
1977 tail = this->tailIndex.load(std::memory_order_acquire);
1978 if ((details::likely)(details::circular_less_than<index_t>(myDequeueCount - overcommit, tail))) {
1979 // Guaranteed to be at least one element to dequeue!
1980
1981 // Get the index. Note that since there's guaranteed to be at least one element, this
1982 // will never exceed tail. We need to do an acquire-release fence here since it's possible
1983 // that whatever condition got us to this point was for an earlier enqueued element (that
1984 // we already see the memory effects for), but that by the time we increment somebody else
1985 // has incremented it, and we need to see the memory effects for *that* element, which is
1986 // in such a case is necessarily visible on the thread that incremented it in the first
1987 // place with the more current condition (they must have acquired a tail that is at least
1988 // as recent).
1989 auto index = this->headIndex.fetch_add(1, std::memory_order_acq_rel);
1990
1991
1992 // Determine which block the element is in
1993
1994 auto localBlockIndex = blockIndex.load(std::memory_order_acquire);
1995 auto localBlockIndexHead = localBlockIndex->front.load(std::memory_order_acquire);
1996
1997 // We need to be careful here about subtracting and dividing because of index wrap-around.
1998 // When an index wraps, we need to preserve the sign of the offset when dividing it by the
1999 // block size (in order to get a correct signed block count offset in all cases):
2000 auto headBase = localBlockIndex->entries[localBlockIndexHead].base;
2001 auto blockBaseIndex = index & ~static_cast<index_t>(BLOCK_SIZE - 1);
2002 auto offset = static_cast<size_t>(static_cast<typename std::make_signed<index_t>::type>(blockBaseIndex - headBase) / BLOCK_SIZE);
2003 auto block = localBlockIndex->entries[(localBlockIndexHead + offset) & (localBlockIndex->size - 1)].block;
2004
2005 // Dequeue
2006 auto& el = *((*block)[index]);
2007 if (!MOODYCAMEL_NOEXCEPT_ASSIGN(T, T&&, element = std::move(el))) {
2008 // Make sure the element is still fully dequeued and destroyed even if the assignment
2009 // throws
2010 struct Guard {
2011 Block* block;
2012 index_t index;
2013
2014 ~Guard()
2015 {
2016 (*block)[index]->~T();
2017 block->ConcurrentQueue::Block::template set_empty<explicit_context>(index);
2018 }
2019 } guard = { block, index };
2020
2021 element = std::move(el); // NOLINT
2022 }
2023 else {
2024 element = std::move(el); // NOLINT
2025 el.~T(); // NOLINT
2026 block->ConcurrentQueue::Block::template set_empty<explicit_context>(index);
2027 }
2028
2029 return true;
2030 }
2031 else {
2032 // Wasn't anything to dequeue after all; make the effective dequeue count eventually consistent
2033 this->dequeueOvercommit.fetch_add(1, std::memory_order_release); // Release so that the fetch_add on dequeueOptimisticCount is guaranteed to happen before this write
2034 }
2035 }
2036
2037 return false;
2038 }
2039
2040 template<AllocationMode allocMode, typename It>
2041 bool MOODYCAMEL_NO_TSAN enqueue_bulk(It itemFirst, size_t count)
2042 {
2043 // First, we need to make sure we have enough room to enqueue all of the elements;
2044 // this means pre-allocating blocks and putting them in the block index (but only if
2045 // all the allocations succeeded).
2046 index_t startTailIndex = this->tailIndex.load(std::memory_order_relaxed);
2047 auto startBlock = this->tailBlock;
2048 auto originalBlockIndexFront = pr_blockIndexFront;
2049 auto originalBlockIndexSlotsUsed = pr_blockIndexSlotsUsed;
2050
2051 Block* firstAllocatedBlock = nullptr;
2052
2053 // Figure out how many blocks we'll need to allocate, and do so
2054 size_t blockBaseDiff = ((startTailIndex + count - 1) & ~static_cast<index_t>(BLOCK_SIZE - 1)) - ((startTailIndex - 1) & ~static_cast<index_t>(BLOCK_SIZE - 1));
2055 index_t currentTailIndex = (startTailIndex - 1) & ~static_cast<index_t>(BLOCK_SIZE - 1);
2056 if (blockBaseDiff > 0) {
2057 // Allocate as many blocks as possible from ahead
2058 while (blockBaseDiff > 0 && this->tailBlock != nullptr && this->tailBlock->next != firstAllocatedBlock && this->tailBlock->next->ConcurrentQueue::Block::template is_empty<explicit_context>()) {
2059 blockBaseDiff -= static_cast<index_t>(BLOCK_SIZE);
2060 currentTailIndex += static_cast<index_t>(BLOCK_SIZE);
2061
2062 this->tailBlock = this->tailBlock->next;
2063 firstAllocatedBlock = firstAllocatedBlock == nullptr ? this->tailBlock : firstAllocatedBlock;
2064
2065 auto& entry = blockIndex.load(std::memory_order_relaxed)->entries[pr_blockIndexFront];
2066 entry.base = currentTailIndex;
2067 entry.block = this->tailBlock;
2068 pr_blockIndexFront = (pr_blockIndexFront + 1) & (pr_blockIndexSize - 1);
2069 }
2070
2071 // Now allocate as many blocks as necessary from the block pool
2072 while (blockBaseDiff > 0) {
2073 blockBaseDiff -= static_cast<index_t>(BLOCK_SIZE);
2074 currentTailIndex += static_cast<index_t>(BLOCK_SIZE);
2075
2076 auto head = this->headIndex.load(std::memory_order_relaxed);
2077 assert(!details::circular_less_than<index_t>(currentTailIndex, head));
2078 bool full = !details::circular_less_than<index_t>(head, currentTailIndex + BLOCK_SIZE) || (MAX_SUBQUEUE_SIZE != details::const_numeric_max<size_t>::value && (MAX_SUBQUEUE_SIZE == 0 || MAX_SUBQUEUE_SIZE - BLOCK_SIZE < currentTailIndex - head));
2079 if (pr_blockIndexRaw == nullptr || pr_blockIndexSlotsUsed == pr_blockIndexSize || full) {
2080 MOODYCAMEL_CONSTEXPR_IF (allocMode == CannotAlloc) {
2081 // Failed to allocate, undo changes (but keep injected blocks)
2082 pr_blockIndexFront = originalBlockIndexFront;
2083 pr_blockIndexSlotsUsed = originalBlockIndexSlotsUsed;
2084 this->tailBlock = startBlock == nullptr ? firstAllocatedBlock : startBlock;
2085 return false;
2086 }
2087 else if (full || !new_block_index(originalBlockIndexSlotsUsed)) {
2088 // Failed to allocate, undo changes (but keep injected blocks)
2089 pr_blockIndexFront = originalBlockIndexFront;
2090 pr_blockIndexSlotsUsed = originalBlockIndexSlotsUsed;
2091 this->tailBlock = startBlock == nullptr ? firstAllocatedBlock : startBlock;
2092 return false;
2093 }
2094
2095 // pr_blockIndexFront is updated inside new_block_index, so we need to
2096 // update our fallback value too (since we keep the new index even if we
2097 // later fail)
2099 }
2100
2101 // Insert a new block in the circular linked list
2102 auto newBlock = this->parent->ConcurrentQueue::template requisition_block<allocMode>();
2103 if (newBlock == nullptr) {
2104 pr_blockIndexFront = originalBlockIndexFront;
2105 pr_blockIndexSlotsUsed = originalBlockIndexSlotsUsed;
2106 this->tailBlock = startBlock == nullptr ? firstAllocatedBlock : startBlock;
2107 return false;
2108 }
2109
2110#ifdef MCDBGQ_TRACKMEM
2111 newBlock->owner = this;
2112#endif
2113 newBlock->ConcurrentQueue::Block::template set_all_empty<explicit_context>();
2114 if (this->tailBlock == nullptr) {
2115 newBlock->next = newBlock;
2116 }
2117 else {
2118 newBlock->next = this->tailBlock->next;
2119 this->tailBlock->next = newBlock;
2120 }
2121 this->tailBlock = newBlock;
2122 firstAllocatedBlock = firstAllocatedBlock == nullptr ? this->tailBlock : firstAllocatedBlock;
2123
2124 ++pr_blockIndexSlotsUsed;
2125
2126 auto& entry = blockIndex.load(std::memory_order_relaxed)->entries[pr_blockIndexFront];
2127 entry.base = currentTailIndex;
2128 entry.block = this->tailBlock;
2129 pr_blockIndexFront = (pr_blockIndexFront + 1) & (pr_blockIndexSize - 1);
2130 }
2131
2132 // Excellent, all allocations succeeded. Reset each block's emptiness before we fill them up, and
2133 // publish the new block index front
2134 auto block = firstAllocatedBlock;
2135 while (true) {
2136 block->ConcurrentQueue::Block::template reset_empty<explicit_context>();
2137 if (block == this->tailBlock) {
2138 break;
2139 }
2140 block = block->next;
2141 }
2142
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);
2145 }
2146 }
2147
2148 // Enqueue, one block at a time
2149 index_t newTailIndex = startTailIndex + static_cast<index_t>(count);
2151 auto endBlock = this->tailBlock;
2152 this->tailBlock = startBlock;
2153 assert((startTailIndex & static_cast<index_t>(BLOCK_SIZE - 1)) != 0 || firstAllocatedBlock != nullptr || count == 0);
2154 if ((startTailIndex & static_cast<index_t>(BLOCK_SIZE - 1)) == 0 && firstAllocatedBlock != nullptr) {
2155 this->tailBlock = firstAllocatedBlock;
2156 }
2157 while (true) {
2158 index_t stopIndex = (currentTailIndex & ~static_cast<index_t>(BLOCK_SIZE - 1)) + static_cast<index_t>(BLOCK_SIZE);
2159 if (details::circular_less_than<index_t>(newTailIndex, stopIndex)) {
2161 }
2162 MOODYCAMEL_CONSTEXPR_IF (MOODYCAMEL_NOEXCEPT_CTOR(T, decltype(*itemFirst), new (static_cast<T*>(nullptr)) T(details::deref_noexcept(itemFirst)))) {
2163 while (currentTailIndex != stopIndex) {
2164 new ((*this->tailBlock)[currentTailIndex++]) T(*itemFirst++);
2165 }
2166 }
2167 else {
2168 MOODYCAMEL_TRY {
2169 while (currentTailIndex != stopIndex) {
2170 // Must use copy constructor even if move constructor is available
2171 // because we may have to revert if there's an exception.
2172 // Sorry about the horrible templated next line, but it was the only way
2173 // to disable moving *at compile time*, which is important because a type
2174 // may only define a (noexcept) move constructor, and so calls to the
2175 // cctor will not compile, even if they are in an if branch that will never
2176 // be executed
2177 new ((*this->tailBlock)[currentTailIndex]) T(details::nomove_if<!MOODYCAMEL_NOEXCEPT_CTOR(T, decltype(*itemFirst), new (static_cast<T*>(nullptr)) T(details::deref_noexcept(itemFirst)))>::eval(*itemFirst));
2179 ++itemFirst;
2180 }
2181 }
2182 MOODYCAMEL_CATCH (...) {
2183 // Oh dear, an exception's been thrown -- destroy the elements that
2184 // were enqueued so far and revert the entire bulk operation (we'll keep
2185 // any allocated blocks in our linked list for later, though).
2187 auto lastBlockEnqueued = this->tailBlock;
2188
2189 pr_blockIndexFront = originalBlockIndexFront;
2190 pr_blockIndexSlotsUsed = originalBlockIndexSlotsUsed;
2191 this->tailBlock = startBlock == nullptr ? firstAllocatedBlock : startBlock;
2192
2194 auto block = startBlock;
2195 if ((startTailIndex & static_cast<index_t>(BLOCK_SIZE - 1)) == 0) {
2196 block = firstAllocatedBlock;
2197 }
2199 while (true) {
2200 stopIndex = (currentTailIndex & ~static_cast<index_t>(BLOCK_SIZE - 1)) + static_cast<index_t>(BLOCK_SIZE);
2201 if (details::circular_less_than<index_t>(constructedStopIndex, stopIndex)) {
2203 }
2204 while (currentTailIndex != stopIndex) {
2205 (*block)[currentTailIndex++]->~T();
2206 }
2207 if (block == lastBlockEnqueued) {
2208 break;
2209 }
2210 block = block->next;
2211 }
2212 }
2213 MOODYCAMEL_RETHROW;
2214 }
2215 }
2216
2217 if (this->tailBlock == endBlock) {
2219 break;
2220 }
2221 this->tailBlock = this->tailBlock->next;
2222 }
2223
2224 MOODYCAMEL_CONSTEXPR_IF (!MOODYCAMEL_NOEXCEPT_CTOR(T, decltype(*itemFirst), new (static_cast<T*>(nullptr)) T(details::deref_noexcept(itemFirst)))) {
2225 if (firstAllocatedBlock != nullptr)
2226 blockIndex.load(std::memory_order_relaxed)->front.store((pr_blockIndexFront - 1) & (pr_blockIndexSize - 1), std::memory_order_release);
2227 }
2228
2229 this->tailIndex.store(newTailIndex, std::memory_order_release);
2230 return true;
2231 }
2232
2233 template<typename It>
2234 size_t dequeue_bulk(It& itemFirst, size_t max)
2235 {
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)) {
2240 desiredCount = desiredCount < max ? desiredCount : max;
2241 std::atomic_thread_fence(std::memory_order_acquire);
2242
2243 auto myDequeueCount = this->dequeueOptimisticCount.fetch_add(desiredCount, std::memory_order_relaxed);
2244
2245 tail = this->tailIndex.load(std::memory_order_acquire);
2246 auto actualCount = static_cast<size_t>(tail - (myDequeueCount - overcommit));
2247 if (details::circular_less_than<size_t>(0, actualCount)) {
2249 if (actualCount < desiredCount) {
2250 this->dequeueOvercommit.fetch_add(desiredCount - actualCount, std::memory_order_release);
2251 }
2252
2253 // Get the first index. Note that since there's guaranteed to be at least actualCount elements, this
2254 // will never exceed tail.
2255 auto firstIndex = this->headIndex.fetch_add(actualCount, std::memory_order_acq_rel);
2256
2257 // Determine which block the first element is in
2258 auto localBlockIndex = blockIndex.load(std::memory_order_acquire);
2259 auto localBlockIndexHead = localBlockIndex->front.load(std::memory_order_acquire);
2260
2261 auto headBase = localBlockIndex->entries[localBlockIndexHead].base;
2262 auto firstBlockBaseIndex = firstIndex & ~static_cast<index_t>(BLOCK_SIZE - 1);
2263 auto offset = static_cast<size_t>(static_cast<typename std::make_signed<index_t>::type>(firstBlockBaseIndex - headBase) / BLOCK_SIZE);
2264 auto indexIndex = (localBlockIndexHead + offset) & (localBlockIndex->size - 1);
2265
2266 // Iterate the blocks and dequeue
2267 auto index = firstIndex;
2268 do {
2269 auto firstIndexInBlock = index;
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;
2272 auto block = localBlockIndex->entries[indexIndex].block;
2273 if (MOODYCAMEL_NOEXCEPT_ASSIGN(T, T&&, details::deref_noexcept(itemFirst) = std::move((*(*block)[index])))) {
2274 while (index != endIndex) {
2275 auto& el = *((*block)[index]);
2276 *itemFirst++ = std::move(el);
2277 el.~T();
2278 ++index;
2279 }
2280 }
2281 else {
2282 MOODYCAMEL_TRY {
2283 while (index != endIndex) {
2284 auto& el = *((*block)[index]);
2285 *itemFirst = std::move(el);
2286 ++itemFirst;
2287 el.~T();
2288 ++index;
2289 }
2290 }
2291 MOODYCAMEL_CATCH (...) {
2292 // It's too late to revert the dequeue, but we can make sure that all
2293 // the dequeued objects are properly destroyed and the block index
2294 // (and empty count) are properly updated before we propagate the exception
2295 do {
2296 block = localBlockIndex->entries[indexIndex].block;
2297 while (index != endIndex) {
2298 (*block)[index++]->~T();
2299 }
2300 block->ConcurrentQueue::Block::template set_many_empty<explicit_context>(firstIndexInBlock, static_cast<size_t>(endIndex - firstIndexInBlock));
2301 indexIndex = (indexIndex + 1) & (localBlockIndex->size - 1);
2302
2303 firstIndexInBlock = index;
2304 endIndex = (index & ~static_cast<index_t>(BLOCK_SIZE - 1)) + static_cast<index_t>(BLOCK_SIZE);
2305 endIndex = details::circular_less_than<index_t>(firstIndex + static_cast<index_t>(actualCount), endIndex) ? firstIndex + static_cast<index_t>(actualCount) : endIndex;
2306 } while (index != firstIndex + actualCount);
2307
2308 MOODYCAMEL_RETHROW;
2309 }
2310 }
2311 block->ConcurrentQueue::Block::template set_many_empty<explicit_context>(firstIndexInBlock, static_cast<size_t>(endIndex - firstIndexInBlock));
2312 indexIndex = (indexIndex + 1) & (localBlockIndex->size - 1);
2313 } while (index != firstIndex + actualCount);
2314
2315 return actualCount;
2316 }
2317 else {
2318 // Wasn't anything to dequeue after all; make the effective dequeue count eventually consistent
2319 this->dequeueOvercommit.fetch_add(desiredCount, std::memory_order_release);
2320 }
2321 }
2322
2323 return 0;
2324 }
2325
2326 private:
2327 struct BlockIndexEntry
2328 {
2329 index_t base;
2330 Block* block;
2331 };
2332
2333 struct BlockIndexHeader
2334 {
2335 size_t size;
2336 std::atomic<size_t> front; // Current slot (not next, like pr_blockIndexFront)
2337 BlockIndexEntry* entries;
2338 void* prev;
2339 };
2340
2341
2342 bool new_block_index(size_t numberOfFilledSlotsToExpose)
2343 {
2344 auto prevBlockSizeMask = pr_blockIndexSize - 1;
2345
2346 // Create the new block
2347 pr_blockIndexSize <<= 1;
2348 auto newRawPtr = static_cast<char*>((Traits::malloc)(sizeof(BlockIndexHeader) + std::alignment_of<BlockIndexEntry>::value - 1 + sizeof(BlockIndexEntry) * pr_blockIndexSize));
2349 if (newRawPtr == nullptr) {
2350 pr_blockIndexSize >>= 1; // Reset to allow graceful retry
2351 return false;
2352 }
2353
2354 auto newBlockIndexEntries = reinterpret_cast<BlockIndexEntry*>(details::align_for<BlockIndexEntry>(newRawPtr + sizeof(BlockIndexHeader)));
2355
2356 // Copy in all the old indices, if any
2357 size_t j = 0;
2358 if (pr_blockIndexSlotsUsed != 0) {
2359 auto i = (pr_blockIndexFront - pr_blockIndexSlotsUsed) & prevBlockSizeMask;
2360 do {
2361 newBlockIndexEntries[j++] = pr_blockIndexEntries[i];
2362 i = (i + 1) & prevBlockSizeMask;
2363 } while (i != pr_blockIndexFront);
2364 }
2365
2366 // Update everything
2367 auto header = new (newRawPtr) BlockIndexHeader;
2368 header->size = pr_blockIndexSize;
2369 header->front.store(numberOfFilledSlotsToExpose - 1, std::memory_order_relaxed);
2370 header->entries = newBlockIndexEntries;
2371 header->prev = pr_blockIndexRaw; // we link the new block to the old one so we can free it later
2372
2373 pr_blockIndexFront = j;
2374 pr_blockIndexEntries = newBlockIndexEntries;
2375 pr_blockIndexRaw = newRawPtr;
2376 blockIndex.store(header, std::memory_order_release);
2377
2378 return true;
2379 }
2380
2381 private:
2382 std::atomic<BlockIndexHeader*> blockIndex;
2383
2384 // To be used by producer only -- consumer must use the ones in referenced by blockIndex
2385 size_t pr_blockIndexSlotsUsed;
2386 size_t pr_blockIndexSize;
2387 size_t pr_blockIndexFront; // Next slot (not current)
2388 BlockIndexEntry* pr_blockIndexEntries;
2389 void* pr_blockIndexRaw;
2390
2391#ifdef MOODYCAMEL_QUEUE_INTERNAL_DEBUG
2392 public:
2393 ExplicitProducer* nextExplicitProducer;
2394 private:
2395#endif
2396
2397#ifdef MCDBGQ_TRACKMEM
2398 friend struct MemStats;
2399#endif
2400 };
2401
2402
2404 // Implicit queue
2406
2407 struct ImplicitProducer : public ProducerBase
2408 {
2409 ImplicitProducer(ConcurrentQueue* parent_) :
2410 ProducerBase(parent_, false),
2411 nextBlockIndexCapacity(IMPLICIT_INITIAL_INDEX_SIZE),
2412 blockIndex(nullptr)
2413 {
2414 new_block_index();
2415 }
2416
2417 ~ImplicitProducer()
2418 {
2419 // Note that since we're in the destructor we can assume that all enqueue/dequeue operations
2420 // completed already; this means that all undequeued elements are placed contiguously across
2421 // contiguous blocks, and that only the first and last remaining blocks can be only partially
2422 // empty (all other remaining blocks must be completely full).
2423
2424#ifdef MOODYCAMEL_CPP11_THREAD_LOCAL_SUPPORTED
2425 // Unregister ourselves for thread termination notification
2426 if (!this->inactive.load(std::memory_order_relaxed)) {
2427 details::ThreadExitNotifier::unsubscribe(&threadExitListener);
2428 }
2429#endif
2430
2431 // Destroy all remaining elements!
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));
2436 bool forceFreeLastBlock = index != tail; // If we enter the loop, then the last (tail) block will not be freed
2437 while (index != tail) {
2438 if ((index & static_cast<index_t>(BLOCK_SIZE - 1)) == 0 || block == nullptr) {
2439 if (block != nullptr) {
2440 // Free the old block
2441 this->parent->add_block_to_free_list(block);
2442 }
2443
2444 block = get_block_index_entry_for_index(index)->value.load(std::memory_order_relaxed);
2445 }
2446
2447 ((*block)[index])->~T();
2448 ++index;
2449 }
2450 // Even if the queue is empty, there's still one block that's not on the free list
2451 // (unless the head index reached the end of it, in which case the tail will be poised
2452 // to create a new block).
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);
2455 }
2456
2457 // Destroy block index
2458 auto localBlockIndex = blockIndex.load(std::memory_order_relaxed);
2459 if (localBlockIndex != nullptr) {
2460 for (size_t i = 0; i != localBlockIndex->capacity; ++i) {
2461 localBlockIndex->index[i]->~BlockIndexEntry();
2462 }
2463 do {
2464 auto prev = localBlockIndex->prev;
2465 localBlockIndex->~BlockIndexHeader();
2466 (Traits::free)(localBlockIndex);
2467 localBlockIndex = prev;
2468 } while (localBlockIndex != nullptr);
2469 }
2470 }
2471
2472 template<AllocationMode allocMode, typename U>
2473 inline bool enqueue(U&& element)
2474 {
2475 index_t currentTailIndex = this->tailIndex.load(std::memory_order_relaxed);
2476 index_t newTailIndex = 1 + currentTailIndex;
2477 if ((currentTailIndex & static_cast<index_t>(BLOCK_SIZE - 1)) == 0) {
2478 // We reached the end of a block, start a new one
2479 auto head = this->headIndex.load(std::memory_order_relaxed);
2480 assert(!details::circular_less_than<index_t>(currentTailIndex, head));
2481 if (!details::circular_less_than<index_t>(head, currentTailIndex + BLOCK_SIZE) || (MAX_SUBQUEUE_SIZE != details::const_numeric_max<size_t>::value && (MAX_SUBQUEUE_SIZE == 0 || MAX_SUBQUEUE_SIZE - BLOCK_SIZE < currentTailIndex - head))) {
2482 return false;
2483 }
2484#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODBLOCKINDEX
2485 debug::DebugLock lock(mutex);
2486#endif
2487 // Find out where we'll be inserting this block in the block index
2488 BlockIndexEntry* idxEntry;
2490 return false;
2491 }
2492
2493 // Get ahold of a new block
2494 auto newBlock = this->parent->ConcurrentQueue::template requisition_block<allocMode>();
2495 if (newBlock == nullptr) {
2496 rewind_block_index_tail();
2497 idxEntry->value.store(nullptr, std::memory_order_relaxed);
2498 return false;
2499 }
2500#ifdef MCDBGQ_TRACKMEM
2501 newBlock->owner = this;
2502#endif
2503 newBlock->ConcurrentQueue::Block::template reset_empty<implicit_context>();
2504
2505 MOODYCAMEL_CONSTEXPR_IF (!MOODYCAMEL_NOEXCEPT_CTOR(T, U, new (static_cast<T*>(nullptr)) T(std::forward<U>(element)))) {
2506 // May throw, try to insert now before we publish the fact that we have this new block
2507 MOODYCAMEL_TRY {
2508 new ((*newBlock)[currentTailIndex]) T(std::forward<U>(element));
2509 }
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);
2514 MOODYCAMEL_RETHROW;
2515 }
2516 }
2517
2518 // Insert the new block into the index
2519 idxEntry->value.store(newBlock, std::memory_order_relaxed);
2520
2521 this->tailBlock = newBlock;
2522
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);
2525 return true;
2526 }
2527 }
2528
2529 // Enqueue
2530 new ((*this->tailBlock)[currentTailIndex]) T(std::forward<U>(element));
2531
2532 this->tailIndex.store(newTailIndex, std::memory_order_release);
2533 return true;
2534 }
2535
2536 template<typename U>
2537 bool dequeue(U& element)
2538 {
2539 // See ExplicitProducer::dequeue for rationale and explanation
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);
2544
2545 index_t myDequeueCount = this->dequeueOptimisticCount.fetch_add(1, std::memory_order_relaxed);
2546 tail = this->tailIndex.load(std::memory_order_acquire);
2547 if ((details::likely)(details::circular_less_than<index_t>(myDequeueCount - overcommit, tail))) {
2548 index_t index = this->headIndex.fetch_add(1, std::memory_order_acq_rel);
2549
2550 // Determine which block the element is in
2551 auto entry = get_block_index_entry_for_index(index);
2552
2553 // Dequeue
2554 auto block = entry->value.load(std::memory_order_relaxed);
2555 auto& el = *((*block)[index]);
2556
2557 if (!MOODYCAMEL_NOEXCEPT_ASSIGN(T, T&&, element = std::move(el))) {
2558#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODBLOCKINDEX
2559 // Note: Acquiring the mutex with every dequeue instead of only when a block
2560 // is released is very sub-optimal, but it is, after all, purely debug code.
2561 debug::DebugLock lock(producer->mutex);
2562#endif
2563 struct Guard {
2564 Block* block;
2565 index_t index;
2566 BlockIndexEntry* entry;
2567 ConcurrentQueue* parent;
2568
2569 ~Guard()
2570 {
2571 (*block)[index]->~T();
2572 if (block->ConcurrentQueue::Block::template set_empty<implicit_context>(index)) {
2573 entry->value.store(nullptr, std::memory_order_relaxed);
2574 parent->add_block_to_free_list(block);
2575 }
2576 }
2577 } guard = { block, index, entry, this->parent };
2578
2579 element = std::move(el); // NOLINT
2580 }
2581 else {
2582 element = std::move(el); // NOLINT
2583 el.~T(); // NOLINT
2584
2585 if (block->ConcurrentQueue::Block::template set_empty<implicit_context>(index)) {
2586 {
2587#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODBLOCKINDEX
2588 debug::DebugLock lock(mutex);
2589#endif
2590 // Add the block back into the global free pool (and remove from block index)
2591 entry->value.store(nullptr, std::memory_order_relaxed);
2592 }
2593 this->parent->add_block_to_free_list(block); // releases the above store
2594 }
2595 }
2596
2597 return true;
2598 }
2599 else {
2600 this->dequeueOvercommit.fetch_add(1, std::memory_order_release);
2601 }
2602 }
2603
2604 return false;
2605 }
2606
2607#ifdef _MSC_VER
2608#pragma warning(push)
2609#pragma warning(disable: 4706) // assignment within conditional expression
2610#endif
2611 template<AllocationMode allocMode, typename It>
2612 bool enqueue_bulk(It itemFirst, size_t count)
2613 {
2614 // First, we need to make sure we have enough room to enqueue all of the elements;
2615 // this means pre-allocating blocks and putting them in the block index (but only if
2616 // all the allocations succeeded).
2617
2618 // Note that the tailBlock we start off with may not be owned by us any more;
2619 // this happens if it was filled up exactly to the top (setting tailIndex to
2620 // the first index of the next block which is not yet allocated), then dequeued
2621 // completely (putting it on the free list) before we enqueue again.
2622
2623 index_t startTailIndex = this->tailIndex.load(std::memory_order_relaxed);
2624 auto startBlock = this->tailBlock;
2625 Block* firstAllocatedBlock = nullptr;
2626 auto endBlock = this->tailBlock;
2627
2628 // Figure out how many blocks we'll need to allocate, and do so
2629 size_t blockBaseDiff = ((startTailIndex + count - 1) & ~static_cast<index_t>(BLOCK_SIZE - 1)) - ((startTailIndex - 1) & ~static_cast<index_t>(BLOCK_SIZE - 1));
2630 index_t currentTailIndex = (startTailIndex - 1) & ~static_cast<index_t>(BLOCK_SIZE - 1);
2631 if (blockBaseDiff > 0) {
2632#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODBLOCKINDEX
2633 debug::DebugLock lock(mutex);
2634#endif
2635 do {
2636 blockBaseDiff -= static_cast<index_t>(BLOCK_SIZE);
2637 currentTailIndex += static_cast<index_t>(BLOCK_SIZE);
2638
2639 // Find out where we'll be inserting this block in the block index
2640 BlockIndexEntry* idxEntry = nullptr; // initialization here unnecessary but compiler can't always tell
2641 Block* newBlock;
2642 bool indexInserted = false;
2643 auto head = this->headIndex.load(std::memory_order_relaxed);
2644 assert(!details::circular_less_than<index_t>(currentTailIndex, head));
2645 bool full = !details::circular_less_than<index_t>(head, currentTailIndex + BLOCK_SIZE) || (MAX_SUBQUEUE_SIZE != details::const_numeric_max<size_t>::value && (MAX_SUBQUEUE_SIZE == 0 || MAX_SUBQUEUE_SIZE - BLOCK_SIZE < currentTailIndex - head));
2646
2647 if (full || !(indexInserted = insert_block_index_entry<allocMode>(idxEntry, currentTailIndex)) || (newBlock = this->parent->ConcurrentQueue::template requisition_block<allocMode>()) == nullptr) {
2648 // Index allocation or block allocation failed; revert any other allocations
2649 // and index insertions done so far for this operation
2650 if (indexInserted) {
2651 rewind_block_index_tail();
2652 idxEntry->value.store(nullptr, std::memory_order_relaxed);
2653 }
2654 currentTailIndex = (startTailIndex - 1) & ~static_cast<index_t>(BLOCK_SIZE - 1);
2655 for (auto block = firstAllocatedBlock; block != nullptr; block = block->next) {
2656 currentTailIndex += static_cast<index_t>(BLOCK_SIZE);
2657 idxEntry = get_block_index_entry_for_index(currentTailIndex);
2658 idxEntry->value.store(nullptr, std::memory_order_relaxed);
2659 rewind_block_index_tail();
2660 }
2661 this->parent->add_blocks_to_free_list(firstAllocatedBlock);
2662 this->tailBlock = startBlock;
2663
2664 return false;
2665 }
2666
2667#ifdef MCDBGQ_TRACKMEM
2668 newBlock->owner = this;
2669#endif
2670 newBlock->ConcurrentQueue::Block::template reset_empty<implicit_context>();
2671 newBlock->next = nullptr;
2672
2673 // Insert the new block into the index
2674 idxEntry->value.store(newBlock, std::memory_order_relaxed);
2675
2676 // Store the chain of blocks so that we can undo if later allocations fail,
2677 // and so that we can find the blocks when we do the actual enqueueing
2678 if ((startTailIndex & static_cast<index_t>(BLOCK_SIZE - 1)) != 0 || firstAllocatedBlock != nullptr) {
2679 assert(this->tailBlock != nullptr);
2680 this->tailBlock->next = newBlock;
2681 }
2682 this->tailBlock = newBlock;
2685 } while (blockBaseDiff > 0);
2686 }
2687
2688 // Enqueue, one block at a time
2689 index_t newTailIndex = startTailIndex + static_cast<index_t>(count);
2691 this->tailBlock = startBlock;
2692 assert((startTailIndex & static_cast<index_t>(BLOCK_SIZE - 1)) != 0 || firstAllocatedBlock != nullptr || count == 0);
2693 if ((startTailIndex & static_cast<index_t>(BLOCK_SIZE - 1)) == 0 && firstAllocatedBlock != nullptr) {
2694 this->tailBlock = firstAllocatedBlock;
2695 }
2696 while (true) {
2697 index_t stopIndex = (currentTailIndex & ~static_cast<index_t>(BLOCK_SIZE - 1)) + static_cast<index_t>(BLOCK_SIZE);
2698 if (details::circular_less_than<index_t>(newTailIndex, stopIndex)) {
2700 }
2701 MOODYCAMEL_CONSTEXPR_IF (MOODYCAMEL_NOEXCEPT_CTOR(T, decltype(*itemFirst), new (static_cast<T*>(nullptr)) T(details::deref_noexcept(itemFirst)))) {
2702 while (currentTailIndex != stopIndex) {
2703 new ((*this->tailBlock)[currentTailIndex++]) T(*itemFirst++);
2704 }
2705 }
2706 else {
2707 MOODYCAMEL_TRY {
2708 while (currentTailIndex != stopIndex) {
2709 new ((*this->tailBlock)[currentTailIndex]) T(details::nomove_if<!MOODYCAMEL_NOEXCEPT_CTOR(T, decltype(*itemFirst), new (static_cast<T*>(nullptr)) T(details::deref_noexcept(itemFirst)))>::eval(*itemFirst));
2711 ++itemFirst;
2712 }
2713 }
2714 MOODYCAMEL_CATCH (...) {
2716 auto lastBlockEnqueued = this->tailBlock;
2717
2719 auto block = startBlock;
2720 if ((startTailIndex & static_cast<index_t>(BLOCK_SIZE - 1)) == 0) {
2721 block = firstAllocatedBlock;
2722 }
2724 while (true) {
2725 stopIndex = (currentTailIndex & ~static_cast<index_t>(BLOCK_SIZE - 1)) + static_cast<index_t>(BLOCK_SIZE);
2726 if (details::circular_less_than<index_t>(constructedStopIndex, stopIndex)) {
2728 }
2729 while (currentTailIndex != stopIndex) {
2730 (*block)[currentTailIndex++]->~T();
2731 }
2732 if (block == lastBlockEnqueued) {
2733 break;
2734 }
2735 block = block->next;
2736 }
2737 }
2738
2739 currentTailIndex = (startTailIndex - 1) & ~static_cast<index_t>(BLOCK_SIZE - 1);
2740 for (auto block = firstAllocatedBlock; block != nullptr; block = block->next) {
2741 currentTailIndex += static_cast<index_t>(BLOCK_SIZE);
2742 auto idxEntry = get_block_index_entry_for_index(currentTailIndex);
2743 idxEntry->value.store(nullptr, std::memory_order_relaxed);
2744 rewind_block_index_tail();
2745 }
2746 this->parent->add_blocks_to_free_list(firstAllocatedBlock);
2747 this->tailBlock = startBlock;
2748 MOODYCAMEL_RETHROW;
2749 }
2750 }
2751
2752 if (this->tailBlock == endBlock) {
2754 break;
2755 }
2756 this->tailBlock = this->tailBlock->next;
2757 }
2758 this->tailIndex.store(newTailIndex, std::memory_order_release);
2759 return true;
2760 }
2761#ifdef _MSC_VER
2762#pragma warning(pop)
2763#endif
2764
2765 template<typename It>
2766 size_t dequeue_bulk(It& itemFirst, size_t max)
2767 {
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)) {
2772 desiredCount = desiredCount < max ? desiredCount : max;
2773 std::atomic_thread_fence(std::memory_order_acquire);
2774
2775 auto myDequeueCount = this->dequeueOptimisticCount.fetch_add(desiredCount, std::memory_order_relaxed);
2776
2777 tail = this->tailIndex.load(std::memory_order_acquire);
2778 auto actualCount = static_cast<size_t>(tail - (myDequeueCount - overcommit));
2779 if (details::circular_less_than<size_t>(0, actualCount)) {
2781 if (actualCount < desiredCount) {
2782 this->dequeueOvercommit.fetch_add(desiredCount - actualCount, std::memory_order_release);
2783 }
2784
2785 // Get the first index. Note that since there's guaranteed to be at least actualCount elements, this
2786 // will never exceed tail.
2787 auto firstIndex = this->headIndex.fetch_add(actualCount, std::memory_order_acq_rel);
2788
2789 // Iterate the blocks and dequeue
2790 auto index = firstIndex;
2791 BlockIndexHeader* localBlockIndex;
2792 auto indexIndex = get_block_index_index_for_index(index, localBlockIndex);
2793 do {
2794 auto blockStartIndex = index;
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;
2797
2798 auto entry = localBlockIndex->index[indexIndex];
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]);
2803 *itemFirst++ = std::move(el);
2804 el.~T();
2805 ++index;
2806 }
2807 }
2808 else {
2809 MOODYCAMEL_TRY {
2810 while (index != endIndex) {
2811 auto& el = *((*block)[index]);
2812 *itemFirst = std::move(el);
2813 ++itemFirst;
2814 el.~T();
2815 ++index;
2816 }
2817 }
2818 MOODYCAMEL_CATCH (...) {
2819 do {
2820 entry = localBlockIndex->index[indexIndex];
2821 block = entry->value.load(std::memory_order_relaxed);
2822 while (index != endIndex) {
2823 (*block)[index++]->~T();
2824 }
2825
2826 if (block->ConcurrentQueue::Block::template set_many_empty<implicit_context>(blockStartIndex, static_cast<size_t>(endIndex - blockStartIndex))) {
2827#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODBLOCKINDEX
2828 debug::DebugLock lock(mutex);
2829#endif
2830 entry->value.store(nullptr, std::memory_order_relaxed);
2831 this->parent->add_block_to_free_list(block);
2832 }
2833 indexIndex = (indexIndex + 1) & (localBlockIndex->capacity - 1);
2834
2835 blockStartIndex = index;
2836 endIndex = (index & ~static_cast<index_t>(BLOCK_SIZE - 1)) + static_cast<index_t>(BLOCK_SIZE);
2837 endIndex = details::circular_less_than<index_t>(firstIndex + static_cast<index_t>(actualCount), endIndex) ? firstIndex + static_cast<index_t>(actualCount) : endIndex;
2838 } while (index != firstIndex + actualCount);
2839
2840 MOODYCAMEL_RETHROW;
2841 }
2842 }
2843 if (block->ConcurrentQueue::Block::template set_many_empty<implicit_context>(blockStartIndex, static_cast<size_t>(endIndex - blockStartIndex))) {
2844 {
2845#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODBLOCKINDEX
2846 debug::DebugLock lock(mutex);
2847#endif
2848 // Note that the set_many_empty above did a release, meaning that anybody who acquires the block
2849 // we're about to free can use it safely since our writes (and reads!) will have happened-before then.
2850 entry->value.store(nullptr, std::memory_order_relaxed);
2851 }
2852 this->parent->add_block_to_free_list(block); // releases the above store
2853 }
2854 indexIndex = (indexIndex + 1) & (localBlockIndex->capacity - 1);
2855 } while (index != firstIndex + actualCount);
2856
2857 return actualCount;
2858 }
2859 else {
2860 this->dequeueOvercommit.fetch_add(desiredCount, std::memory_order_release);
2861 }
2862 }
2863
2864 return 0;
2865 }
2866
2867 private:
2868 // The block size must be > 1, so any number with the low bit set is an invalid block base index
2869 static const index_t INVALID_BLOCK_BASE = 1;
2870
2871 struct BlockIndexEntry
2872 {
2873 std::atomic<index_t> key;
2874 std::atomic<Block*> value;
2875 };
2876
2877 struct BlockIndexHeader
2878 {
2879 size_t capacity;
2880 std::atomic<size_t> tail;
2881 BlockIndexEntry* entries;
2882 BlockIndexEntry** index;
2883 BlockIndexHeader* prev;
2884 };
2885
2886 template<AllocationMode allocMode>
2887 inline bool insert_block_index_entry(BlockIndexEntry*& idxEntry, index_t blockStartIndex)
2888 {
2889 auto localBlockIndex = blockIndex.load(std::memory_order_relaxed); // We're the only writer thread, relaxed is OK
2890 if (localBlockIndex == nullptr) {
2891 return false; // this can happen if new_block_index failed in the constructor
2892 }
2893 size_t newTail = (localBlockIndex->tail.load(std::memory_order_relaxed) + 1) & (localBlockIndex->capacity - 1);
2895 if (idxEntry->key.load(std::memory_order_relaxed) == INVALID_BLOCK_BASE ||
2896 idxEntry->value.load(std::memory_order_relaxed) == nullptr) {
2897
2898 idxEntry->key.store(blockStartIndex, std::memory_order_relaxed);
2899 localBlockIndex->tail.store(newTail, std::memory_order_release);
2900 return true;
2901 }
2902
2903 // No room in the old block index, try to allocate another one!
2904 MOODYCAMEL_CONSTEXPR_IF (allocMode == CannotAlloc) {
2905 return false;
2906 }
2907 else if (!new_block_index()) {
2908 return false;
2909 }
2910 localBlockIndex = blockIndex.load(std::memory_order_relaxed);
2911 newTail = (localBlockIndex->tail.load(std::memory_order_relaxed) + 1) & (localBlockIndex->capacity - 1);
2913 assert(idxEntry->key.load(std::memory_order_relaxed) == INVALID_BLOCK_BASE);
2914 idxEntry->key.store(blockStartIndex, std::memory_order_relaxed);
2915 localBlockIndex->tail.store(newTail, std::memory_order_release);
2916 return true;
2917 }
2918
2919 inline void rewind_block_index_tail()
2920 {
2921 auto localBlockIndex = blockIndex.load(std::memory_order_relaxed);
2922 localBlockIndex->tail.store((localBlockIndex->tail.load(std::memory_order_relaxed) - 1) & (localBlockIndex->capacity - 1), std::memory_order_relaxed);
2923 }
2924
2925 inline BlockIndexEntry* get_block_index_entry_for_index(index_t index) const
2926 {
2927 BlockIndexHeader* localBlockIndex;
2928 auto idx = get_block_index_index_for_index(index, localBlockIndex);
2929 return localBlockIndex->index[idx];
2930 }
2931
2932 inline size_t get_block_index_index_for_index(index_t index, BlockIndexHeader*& localBlockIndex) const
2933 {
2934#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODBLOCKINDEX
2935 debug::DebugLock lock(mutex);
2936#endif
2937 index &= ~static_cast<index_t>(BLOCK_SIZE - 1);
2938 localBlockIndex = blockIndex.load(std::memory_order_acquire);
2939 auto tail = localBlockIndex->tail.load(std::memory_order_acquire);
2940 auto tailBase = localBlockIndex->index[tail]->key.load(std::memory_order_relaxed);
2941 assert(tailBase != INVALID_BLOCK_BASE);
2942 // Note: Must use division instead of shift because the index may wrap around, causing a negative
2943 // offset, whose negativity we want to preserve
2944 auto offset = static_cast<size_t>(static_cast<typename std::make_signed<index_t>::type>(index - tailBase) / BLOCK_SIZE);
2945 size_t idx = (tail + offset) & (localBlockIndex->capacity - 1);
2946 assert(localBlockIndex->index[idx]->key.load(std::memory_order_relaxed) == index && localBlockIndex->index[idx]->value.load(std::memory_order_relaxed) != nullptr);
2947 return idx;
2948 }
2949
2950 bool new_block_index()
2951 {
2952 auto prev = blockIndex.load(std::memory_order_relaxed);
2953 size_t prevCapacity = prev == nullptr ? 0 : prev->capacity;
2954 auto entryCount = prev == nullptr ? nextBlockIndexCapacity : prevCapacity;
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) {
2960 return false;
2961 }
2962
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);
2968 auto prevPos = prevTail;
2969 size_t i = 0;
2970 do {
2971 prevPos = (prevPos + 1) & (prev->capacity - 1);
2972 index[i++] = prev->index[prevPos];
2973 } while (prevPos != prevTail);
2974 assert(i == prevCapacity);
2975 }
2976 for (size_t i = 0; i != entryCount; ++i) {
2977 new (entries + i) BlockIndexEntry;
2978 entries[i].key.store(INVALID_BLOCK_BASE, std::memory_order_relaxed);
2979 index[prevCapacity + i] = entries + i;
2980 }
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);
2986
2987 blockIndex.store(header, std::memory_order_release);
2988
2989 nextBlockIndexCapacity <<= 1;
2990
2991 return true;
2992 }
2993
2994 private:
2995 size_t nextBlockIndexCapacity;
2996 std::atomic<BlockIndexHeader*> blockIndex;
2997
2998#ifdef MOODYCAMEL_CPP11_THREAD_LOCAL_SUPPORTED
2999 public:
3000 details::ThreadExitListener threadExitListener;
3001 private:
3002#endif
3003
3004#ifdef MOODYCAMEL_QUEUE_INTERNAL_DEBUG
3005 public:
3006 ImplicitProducer* nextImplicitProducer;
3007 private:
3008#endif
3009
3010#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODBLOCKINDEX
3011 mutable debug::DebugMutex mutex;
3012#endif
3013#ifdef MCDBGQ_TRACKMEM
3014 friend struct MemStats;
3015#endif
3016 };
3017
3018
3020 // Block pool manipulation
3022
3023 void populate_initial_block_list(size_t blockCount)
3024 {
3025 initialBlockPoolSize = blockCount;
3026 if (initialBlockPoolSize == 0) {
3027 initialBlockPool = nullptr;
3028 return;
3029 }
3030
3031 initialBlockPool = create_array<Block>(blockCount);
3032 if (initialBlockPool == nullptr) {
3033 initialBlockPoolSize = 0;
3034 }
3035 for (size_t i = 0; i < initialBlockPoolSize; ++i) {
3036 initialBlockPool[i].dynamicallyAllocated = false;
3037 }
3038 }
3039
3040 inline Block* try_get_block_from_initial_pool()
3041 {
3042 if (initialBlockPoolIndex.load(std::memory_order_relaxed) >= initialBlockPoolSize) {
3043 return nullptr;
3044 }
3045
3046 auto index = initialBlockPoolIndex.fetch_add(1, std::memory_order_relaxed);
3047
3048 return index < initialBlockPoolSize ? (initialBlockPool + index) : nullptr;
3049 }
3050
3051 inline void add_block_to_free_list(Block* block)
3052 {
3053#ifdef MCDBGQ_TRACKMEM
3054 block->owner = nullptr;
3055#endif
3056 freeList.add(block);
3057 }
3058
3059 inline void add_blocks_to_free_list(Block* block)
3060 {
3061 while (block != nullptr) {
3062 auto next = block->next;
3063 add_block_to_free_list(block);
3064 block = next;
3065 }
3066 }
3067
3068 inline Block* try_get_block_from_free_list()
3069 {
3070 return freeList.try_get();
3071 }
3072
3073 // Gets a free block from one of the memory pools, or allocates a new one (if applicable)
3074 template<AllocationMode canAlloc>
3075 Block* requisition_block()
3076 {
3077 auto block = try_get_block_from_initial_pool();
3078 if (block != nullptr) {
3079 return block;
3080 }
3081
3082 block = try_get_block_from_free_list();
3083 if (block != nullptr) {
3084 return block;
3085 }
3086
3087 MOODYCAMEL_CONSTEXPR_IF (canAlloc == CanAlloc) {
3088 return create<Block>();
3089 }
3090 else {
3091 return nullptr;
3092 }
3093 }
3094
3095
3096#ifdef MCDBGQ_TRACKMEM
3097 public:
3098 struct MemStats {
3099 size_t allocatedBlocks;
3100 size_t usedBlocks;
3101 size_t freeBlocks;
3102 size_t ownedBlocksExplicit;
3103 size_t ownedBlocksImplicit;
3104 size_t implicitProducers;
3105 size_t explicitProducers;
3106 size_t elementsEnqueued;
3107 size_t blockClassBytes;
3108 size_t queueClassBytes;
3111
3112 friend class ConcurrentQueue;
3113
3114 private:
3116 {
3117 MemStats stats = { 0 };
3118
3119 stats.elementsEnqueued = q->size_approx();
3120
3121 auto block = q->freeList.head_unsafe();
3122 while (block != nullptr) {
3123 ++stats.allocatedBlocks;
3124 ++stats.freeBlocks;
3125 block = block->freeListNext.load(std::memory_order_relaxed);
3126 }
3127
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;
3132
3133 if (implicit) {
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;
3144 }
3145 }
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*);
3149 }
3150 }
3151 for (; details::circular_less_than<index_t>(head, tail); head += BLOCK_SIZE) {
3152 //auto block = prod->get_block_index_entry_for_index(head);
3153 ++stats.usedBlocks;
3154 }
3155 }
3156 else {
3157 auto prod = static_cast<ExplicitProducer*>(ptr);
3158 stats.queueClassBytes += sizeof(ExplicitProducer);
3159 auto tailBlock = prod->tailBlock;
3160 bool wasNonEmpty = false;
3161 if (tailBlock != nullptr) {
3162 auto block = tailBlock;
3163 do {
3164 ++stats.allocatedBlocks;
3165 if (!block->ConcurrentQueue::Block::template is_empty<explicit_context>() || wasNonEmpty) {
3166 ++stats.usedBlocks;
3167 wasNonEmpty = wasNonEmpty || block != tailBlock;
3168 }
3169 ++stats.ownedBlocksExplicit;
3170 block = block->next;
3171 } while (block != tailBlock);
3172 }
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);
3177 }
3178 }
3179 }
3180
3181 auto freeOnInitialPool = q->initialBlockPoolIndex.load(std::memory_order_relaxed) >= q->initialBlockPoolSize ? 0 : q->initialBlockPoolSize - q->initialBlockPoolIndex.load(std::memory_order_relaxed);
3182 stats.allocatedBlocks += freeOnInitialPool;
3183 stats.freeBlocks += freeOnInitialPool;
3184
3185 stats.blockClassBytes = sizeof(Block) * stats.allocatedBlocks;
3186 stats.queueClassBytes += sizeof(ConcurrentQueue);
3187
3188 return stats;
3189 }
3190 };
3191
3192 // For debugging only. Not thread-safe.
3194 {
3195 return MemStats::getFor(this);
3196 }
3197 private:
3198 friend struct MemStats;
3199#endif
3200
3201
3203 // Producer list manipulation
3205
3206 ProducerBase* recycle_or_create_producer(bool isExplicit)
3207 {
3208 bool recycled;
3209 return recycle_or_create_producer(isExplicit, recycled);
3210 }
3211
3212 ProducerBase* recycle_or_create_producer(bool isExplicit, bool& recycled)
3213 {
3214#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODHASH
3215 debug::DebugLock lock(implicitProdMutex);
3216#endif
3217 // Try to re-use one first
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, /* desired */ false, std::memory_order_acquire, std::memory_order_relaxed)) {
3222 // We caught one! It's been marked as activated, the caller can have it
3223 recycled = true;
3224 return ptr;
3225 }
3226 }
3227 }
3228
3229 recycled = false;
3230 return add_producer(isExplicit ? static_cast<ProducerBase*>(create<ExplicitProducer>(this)) : create<ImplicitProducer>(this));
3231 }
3232
3233 ProducerBase* add_producer(ProducerBase* producer)
3234 {
3235 // Handle failed memory allocation
3236 if (producer == nullptr) {
3237 return nullptr;
3238 }
3239
3240 producerCount.fetch_add(1, std::memory_order_relaxed);
3241
3242 // Add it to the lock-free list
3243 auto prevTail = producerListTail.load(std::memory_order_relaxed);
3244 do {
3245 producer->next = prevTail;
3246 } while (!producerListTail.compare_exchange_weak(prevTail, producer, std::memory_order_release, std::memory_order_relaxed));
3247
3248#ifdef MOODYCAMEL_QUEUE_INTERNAL_DEBUG
3249 if (producer->isExplicit) {
3250 auto prevTailExplicit = explicitProducers.load(std::memory_order_relaxed);
3251 do {
3252 static_cast<ExplicitProducer*>(producer)->nextExplicitProducer = prevTailExplicit;
3253 } while (!explicitProducers.compare_exchange_weak(prevTailExplicit, static_cast<ExplicitProducer*>(producer), std::memory_order_release, std::memory_order_relaxed));
3254 }
3255 else {
3256 auto prevTailImplicit = implicitProducers.load(std::memory_order_relaxed);
3257 do {
3258 static_cast<ImplicitProducer*>(producer)->nextImplicitProducer = prevTailImplicit;
3259 } while (!implicitProducers.compare_exchange_weak(prevTailImplicit, static_cast<ImplicitProducer*>(producer), std::memory_order_release, std::memory_order_relaxed));
3260 }
3261#endif
3262
3263 return producer;
3264 }
3265
3266 void reown_producers()
3267 {
3268 // After another instance is moved-into/swapped-with this one, all the
3269 // producers we stole still think their parents are the other queue.
3270 // So fix them up!
3271 for (auto ptr = producerListTail.load(std::memory_order_relaxed); ptr != nullptr; ptr = ptr->next_prod()) {
3272 ptr->parent = this;
3273 }
3274 }
3275
3276
3278 // Implicit producer hash
3280
3281 struct ImplicitProducerKVP
3282 {
3283 std::atomic<details::thread_id_t> key;
3284 ImplicitProducer* value; // No need for atomicity since it's only read by the thread that sets it in the first place
3285
3286 ImplicitProducerKVP() : value(nullptr) { }
3287
3288 ImplicitProducerKVP(ImplicitProducerKVP&& other) MOODYCAMEL_NOEXCEPT
3289 {
3290 key.store(other.key.load(std::memory_order_relaxed), std::memory_order_relaxed);
3291 value = other.value;
3292 }
3293
3294 inline ImplicitProducerKVP& operator=(ImplicitProducerKVP&& other) MOODYCAMEL_NOEXCEPT
3295 {
3296 swap(other);
3297 return *this;
3298 }
3299
3300 inline void swap(ImplicitProducerKVP& other) MOODYCAMEL_NOEXCEPT
3301 {
3302 if (this != &other) {
3303 details::swap_relaxed(key, other.key);
3304 std::swap(value, other.value);
3305 }
3306 }
3307 };
3308
3309 template<typename XT, typename XTraits>
3310 friend void moodycamel::swap(typename ConcurrentQueue<XT, XTraits>::ImplicitProducerKVP&, typename ConcurrentQueue<XT, XTraits>::ImplicitProducerKVP&) MOODYCAMEL_NOEXCEPT;
3311
3312 struct ImplicitProducerHash
3313 {
3314 size_t capacity;
3315 ImplicitProducerKVP* entries;
3316 ImplicitProducerHash* prev;
3317 };
3318
3319 inline void populate_initial_implicit_producer_hash()
3320 {
3321 MOODYCAMEL_CONSTEXPR_IF (INITIAL_IMPLICIT_PRODUCER_HASH_SIZE == 0) {
3322 return;
3323 }
3324 else {
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);
3331 }
3332 hash->prev = nullptr;
3333 implicitProducerHash.store(hash, std::memory_order_relaxed);
3334 }
3335 }
3336
3337 void swap_implicit_producer_hashes(ConcurrentQueue& other)
3338 {
3339 MOODYCAMEL_CONSTEXPR_IF (INITIAL_IMPLICIT_PRODUCER_HASH_SIZE == 0) {
3340 return;
3341 }
3342 else {
3343 // Swap (assumes our implicit producer hash is initialized)
3344 initialImplicitProducerHashEntries.swap(other.initialImplicitProducerHashEntries);
3345 initialImplicitProducerHash.entries = &initialImplicitProducerHashEntries[0];
3346 other.initialImplicitProducerHash.entries = &other.initialImplicitProducerHashEntries[0];
3347
3348 details::swap_relaxed(implicitProducerHashCount, other.implicitProducerHashCount);
3349
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);
3353 }
3354 else {
3355 ImplicitProducerHash* hash;
3356 for (hash = implicitProducerHash.load(std::memory_order_relaxed); hash->prev != &other.initialImplicitProducerHash; hash = hash->prev) {
3357 continue;
3358 }
3359 hash->prev = &initialImplicitProducerHash;
3360 }
3361 if (other.implicitProducerHash.load(std::memory_order_relaxed) == &initialImplicitProducerHash) {
3362 other.implicitProducerHash.store(&other.initialImplicitProducerHash, std::memory_order_relaxed);
3363 }
3364 else {
3365 ImplicitProducerHash* hash;
3366 for (hash = other.implicitProducerHash.load(std::memory_order_relaxed); hash->prev != &initialImplicitProducerHash; hash = hash->prev) {
3367 continue;
3368 }
3369 hash->prev = &other.initialImplicitProducerHash;
3370 }
3371 }
3372 }
3373
3374 // Only fails (returns nullptr) if memory allocation fails
3375 ImplicitProducer* get_or_add_implicit_producer()
3376 {
3377 // Note that since the data is essentially thread-local (key is thread ID),
3378 // there's a reduced need for fences (memory ordering is already consistent
3379 // for any individual thread), except for the current table itself.
3380
3381 // Start by looking for the thread ID in the current and all previous hash tables.
3382 // If it's not found, it must not be in there yet, since this same thread would
3383 // have added it previously to one of the tables that we traversed.
3384
3385 // Code and algorithm adapted from http://preshing.com/20130605/the-worlds-simplest-lock-free-hash-table
3386
3387#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODHASH
3388 debug::DebugLock lock(implicitProdMutex);
3389#endif
3390
3391 auto id = details::thread_id();
3392 auto hashedId = details::hash_thread_id(id);
3393
3394 auto mainHash = implicitProducerHash.load(std::memory_order_acquire);
3395 assert(mainHash != nullptr); // silence clang-tidy and MSVC warnings (hash cannot be null)
3396 for (auto hash = mainHash; hash != nullptr; hash = hash->prev) {
3397 // Look for the id in this hash
3398 auto index = hashedId;
3399 while (true) { // Not an infinite loop because at least one slot is free in the hash table
3400 index &= hash->capacity - 1;
3401
3402 auto probedKey = hash->entries[index].key.load(std::memory_order_relaxed);
3403 if (probedKey == id) {
3404 // Found it! If we had to search several hashes deep, though, we should lazily add it
3405 // to the current main hash table to avoid the extended search next time.
3406 // Note there's guaranteed to be room in the current hash table since every subsequent
3407 // table implicitly reserves space for all previous tables (there's only one
3408 // implicitProducerHashCount).
3409 auto value = hash->entries[index].value;
3410 if (hash != mainHash) {
3411 index = hashedId;
3412 while (true) {
3413 index &= mainHash->capacity - 1;
3414 probedKey = mainHash->entries[index].key.load(std::memory_order_relaxed);
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)) ||
3419 (probedKey == reusable && mainHash->entries[index].key.compare_exchange_strong(reusable, id, std::memory_order_acquire, std::memory_order_acquire))) {
3420#else
3421 if ((probedKey == empty && mainHash->entries[index].key.compare_exchange_strong(empty, id, std::memory_order_relaxed, std::memory_order_relaxed))) {
3422#endif
3423 mainHash->entries[index].value = value;
3424 break;
3425 }
3426 ++index;
3427 }
3428 }
3429
3430 return value;
3431 }
3432 if (probedKey == details::invalid_thread_id) {
3433 break; // Not in this hash table
3434 }
3435 ++index;
3436 }
3437 }
3438
3439 // Insert!
3440 auto newCount = 1 + implicitProducerHashCount.fetch_add(1, std::memory_order_relaxed);
3441 while (true) {
3442 // NOLINTNEXTLINE(clang-analyzer-core.NullDereference)
3443 if (newCount >= (mainHash->capacity >> 1) && !implicitProducerHashResizeInProgress.test_and_set(std::memory_order_acquire)) {
3444 // We've acquired the resize lock, try to allocate a bigger hash table.
3445 // Note the acquire fence synchronizes with the release fence at the end of this block, and hence when
3446 // we reload implicitProducerHash it must be the most recent version (it only gets changed within this
3447 // locked block).
3448 mainHash = implicitProducerHash.load(std::memory_order_acquire);
3449 if (newCount >= (mainHash->capacity >> 1)) {
3450 auto newCapacity = mainHash->capacity << 1;
3451 while (newCount >= (newCapacity >> 1)) {
3452 newCapacity <<= 1;
3453 }
3454 auto raw = static_cast<char*>((Traits::malloc)(sizeof(ImplicitProducerHash) + std::alignment_of<ImplicitProducerKVP>::value - 1 + sizeof(ImplicitProducerKVP) * newCapacity));
3455 if (raw == nullptr) {
3456 // Allocation failed
3457 implicitProducerHashCount.fetch_sub(1, std::memory_order_relaxed);
3458 implicitProducerHashResizeInProgress.clear(std::memory_order_relaxed);
3459 return nullptr;
3460 }
3461
3462 auto newHash = new (raw) ImplicitProducerHash;
3463 newHash->capacity = static_cast<size_t>(newCapacity);
3464 newHash->entries = reinterpret_cast<ImplicitProducerKVP*>(details::align_for<ImplicitProducerKVP>(raw + sizeof(ImplicitProducerHash)));
3465 for (size_t i = 0; i != newCapacity; ++i) {
3466 new (newHash->entries + i) ImplicitProducerKVP;
3467 newHash->entries[i].key.store(details::invalid_thread_id, std::memory_order_relaxed);
3468 }
3469 newHash->prev = mainHash;
3470 implicitProducerHash.store(newHash, std::memory_order_release);
3471 implicitProducerHashResizeInProgress.clear(std::memory_order_release);
3472 mainHash = newHash;
3473 }
3474 else {
3475 implicitProducerHashResizeInProgress.clear(std::memory_order_release);
3476 }
3477 }
3478
3479 // If it's < three-quarters full, add to the old one anyway so that we don't have to wait for the next table
3480 // to finish being allocated by another thread (and if we just finished allocating above, the condition will
3481 // always be true)
3482 if (newCount < (mainHash->capacity >> 1) + (mainHash->capacity >> 2)) {
3483 bool recycled;
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);
3487 return nullptr;
3488 }
3489 if (recycled) {
3490 implicitProducerHashCount.fetch_sub(1, std::memory_order_relaxed);
3491 }
3492
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);
3497#endif
3498
3499 auto index = hashedId;
3500 while (true) {
3501 index &= mainHash->capacity - 1;
3502 auto probedKey = mainHash->entries[index].key.load(std::memory_order_relaxed);
3503
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)) ||
3508 (probedKey == reusable && mainHash->entries[index].key.compare_exchange_strong(reusable, id, std::memory_order_acquire, std::memory_order_acquire))) {
3509#else
3510 if ((probedKey == empty && mainHash->entries[index].key.compare_exchange_strong(empty, id, std::memory_order_relaxed, std::memory_order_relaxed))) {
3511#endif
3512 mainHash->entries[index].value = producer;
3513 break;
3514 }
3515 ++index;
3516 }
3517 return producer;
3518 }
3519
3520 // Hmm, the old hash is quite full and somebody else is busy allocating a new one.
3521 // We need to wait for the allocating thread to finish (if it succeeds, we add, if not,
3522 // we try to allocate ourselves).
3523 mainHash = implicitProducerHash.load(std::memory_order_acquire);
3524 }
3525 }
3526
3527#ifdef MOODYCAMEL_CPP11_THREAD_LOCAL_SUPPORTED
3528 void implicit_producer_thread_exited(ImplicitProducer* producer)
3529 {
3530 // Remove from thread exit listeners
3531 details::ThreadExitNotifier::unsubscribe(&producer->threadExitListener);
3532
3533 // Remove from hash
3534#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODHASH
3535 debug::DebugLock lock(implicitProdMutex);
3536#endif
3537 auto hash = implicitProducerHash.load(std::memory_order_acquire);
3538 assert(hash != nullptr); // The thread exit listener is only registered if we were added to a hash in the first place
3539 auto id = details::thread_id();
3540 auto hashedId = details::hash_thread_id(id);
3541 details::thread_id_t probedKey;
3542
3543 // We need to traverse all the hashes just in case other threads aren't on the current one yet and are
3544 // trying to add an entry thinking there's a free slot (because they reused a producer)
3545 for (; hash != nullptr; hash = hash->prev) {
3546 auto index = hashedId;
3547 do {
3548 index &= hash->capacity - 1;
3549 probedKey = hash->entries[index].key.load(std::memory_order_relaxed);
3550 if (probedKey == id) {
3551 hash->entries[index].key.store(details::invalid_thread_id2, std::memory_order_release);
3552 break;
3553 }
3554 ++index;
3555 } while (probedKey != details::invalid_thread_id); // Can happen if the hash has changed but we weren't put back in it yet, or if we weren't added to this hash in the first place
3556 }
3557
3558 // Mark the queue as being recyclable
3559 producer->inactive.store(true, std::memory_order_release);
3560 }
3561
3562 static void implicit_producer_thread_exited_callback(void* userData)
3563 {
3564 auto producer = static_cast<ImplicitProducer*>(userData);
3565 auto queue = producer->parent;
3566 queue->implicit_producer_thread_exited(producer);
3567 }
3568#endif
3569
3571 // Utility functions
3573
3574 template<typename TAlign>
3575 static inline void* aligned_malloc(size_t size)
3576 {
3577 MOODYCAMEL_CONSTEXPR_IF (std::alignment_of<TAlign>::value <= std::alignment_of<details::max_align_t>::value)
3578 return (Traits::malloc)(size);
3579 else {
3580 size_t alignment = std::alignment_of<TAlign>::value;
3581 void* raw = (Traits::malloc)(size + alignment - 1 + sizeof(void*));
3582 if (!raw)
3583 return nullptr;
3584 char* ptr = details::align_for<TAlign>(reinterpret_cast<char*>(raw) + sizeof(void*));
3585 *(reinterpret_cast<void**>(ptr) - 1) = raw;
3586 return ptr;
3587 }
3588 }
3589
3590 template<typename TAlign>
3591 static inline void aligned_free(void* ptr)
3592 {
3593 MOODYCAMEL_CONSTEXPR_IF (std::alignment_of<TAlign>::value <= std::alignment_of<details::max_align_t>::value)
3594 return (Traits::free)(ptr);
3595 else
3596 (Traits::free)(ptr ? *(reinterpret_cast<void**>(ptr) - 1) : nullptr);
3597 }
3598
3599 template<typename U>
3600 static inline U* create_array(size_t count)
3601 {
3602 assert(count > 0);
3603 U* p = static_cast<U*>(aligned_malloc<U>(sizeof(U) * count));
3604 if (p == nullptr)
3605 return nullptr;
3606
3607 for (size_t i = 0; i != count; ++i)
3608 new (p + i) U();
3609 return p;
3610 }
3611
3612 template<typename U>
3613 static inline void destroy_array(U* p, size_t count)
3614 {
3615 if (p != nullptr) {
3616 assert(count > 0);
3617 for (size_t i = count; i != 0; )
3618 (p + --i)->~U();
3619 }
3620 aligned_free<U>(p);
3621 }
3622
3623 template<typename U>
3624 static inline U* create()
3625 {
3626 void* p = aligned_malloc<U>(sizeof(U));
3627 return p != nullptr ? new (p) U : nullptr;
3628 }
3629
3630 template<typename U, typename A1>
3631 static inline U* create(A1&& a1)
3632 {
3633 void* p = aligned_malloc<U>(sizeof(U));
3634 return p != nullptr ? new (p) U(std::forward<A1>(a1)) : nullptr;
3635 }
3636
3637 template<typename U>
3638 static inline void destroy(U* p)
3639 {
3640 if (p != nullptr)
3641 p->~U();
3642 aligned_free<U>(p);
3643 }
3644
3645private:
3646 std::atomic<ProducerBase*> producerListTail;
3647 std::atomic<std::uint32_t> producerCount;
3648
3649 std::atomic<size_t> initialBlockPoolIndex;
3650 Block* initialBlockPool;
3651 size_t initialBlockPoolSize;
3652
3653#ifndef MCDBGQ_USEDEBUGFREELIST
3654 FreeList<Block> freeList;
3655#else
3656 debug::DebugFreeList<Block> freeList;
3657#endif
3658
3659 std::atomic<ImplicitProducerHash*> implicitProducerHash;
3660 std::atomic<size_t> implicitProducerHashCount; // Number of slots logically used
3661 ImplicitProducerHash initialImplicitProducerHash;
3662 std::array<ImplicitProducerKVP, INITIAL_IMPLICIT_PRODUCER_HASH_SIZE> initialImplicitProducerHashEntries;
3663 std::atomic_flag implicitProducerHashResizeInProgress;
3664
3665 std::atomic<std::uint32_t> nextExplicitConsumerId;
3666 std::atomic<std::uint32_t> globalExplicitConsumerOffset;
3667
3668#ifdef MCDBGQ_NOLOCKFREE_IMPLICITPRODHASH
3670#endif
3671
3672#ifdef MOODYCAMEL_QUEUE_INTERNAL_DEBUG
3673 std::atomic<ExplicitProducer*> explicitProducers;
3674 std::atomic<ImplicitProducer*> implicitProducers;
3675#endif
3676};
3677
3678
3679template<typename T, typename Traits>
3680ProducerToken::ProducerToken(ConcurrentQueue<T, Traits>& queue)
3681 : producer(queue.recycle_or_create_producer(true))
3682{
3683 if (producer != nullptr) {
3684 producer->token = this;
3685 }
3686}
3687
3688template<typename T, typename Traits>
3689ProducerToken::ProducerToken(BlockingConcurrentQueue<T, Traits>& queue)
3690 : producer(reinterpret_cast<ConcurrentQueue<T, Traits>*>(&queue)->recycle_or_create_producer(true))
3691{
3692 if (producer != nullptr) {
3693 producer->token = this;
3694 }
3695}
3696
3697template<typename T, typename Traits>
3698ConsumerToken::ConsumerToken(ConcurrentQueue<T, Traits>& queue)
3699 : itemsConsumedFromCurrent(0), currentProducer(nullptr), desiredProducer(nullptr)
3700{
3701 initialOffset = queue.nextExplicitConsumerId.fetch_add(1, std::memory_order_release);
3702 lastKnownGlobalOffset = static_cast<std::uint32_t>(-1);
3703}
3704
3705template<typename T, typename Traits>
3706ConsumerToken::ConsumerToken(BlockingConcurrentQueue<T, Traits>& queue)
3707 : itemsConsumedFromCurrent(0), currentProducer(nullptr), desiredProducer(nullptr)
3708{
3709 initialOffset = reinterpret_cast<ConcurrentQueue<T, Traits>*>(&queue)->nextExplicitConsumerId.fetch_add(1, std::memory_order_release);
3710 lastKnownGlobalOffset = static_cast<std::uint32_t>(-1);
3711}
3712
3713template<typename T, typename Traits>
3714inline void swap(ConcurrentQueue<T, Traits>& a, ConcurrentQueue<T, Traits>& b) MOODYCAMEL_NOEXCEPT
3715{
3716 a.swap(b);
3717}
3718
3719inline void swap(ProducerToken& a, ProducerToken& b) MOODYCAMEL_NOEXCEPT
3720{
3721 a.swap(b);
3722}
3723
3724inline void swap(ConsumerToken& a, ConsumerToken& b) MOODYCAMEL_NOEXCEPT
3725{
3726 a.swap(b);
3727}
3728
3729template<typename T, typename Traits>
3730inline void swap(typename ConcurrentQueue<T, Traits>::ImplicitProducerKVP& a, typename ConcurrentQueue<T, Traits>::ImplicitProducerKVP& b) MOODYCAMEL_NOEXCEPT
3731{
3732 a.swap(b);
3733}
3734
3735}
3736
3737#if defined(_MSC_VER) && (!defined(_HAS_CXX17) || !_HAS_CXX17)
3738#pragma warning(pop)
3739#endif
3740
3741#if defined(__GNUC__)
3742#pragma GCC diagnostic pop
3743#endif
Definition base.h:1940
Definition blockingconcurrentqueue.h:26
Definition concurrentqueue.h:747
Definition concurrentqueue.h:321
Definition concurrentqueue.h:696
Definition concurrentqueue.h:631
Definition concurrentqueue_internal_debug.h:27
Definition concurrentqueue.h:435
Definition concurrentqueue.h:292
Definition concurrentqueue.h:459
Definition concurrentqueue.h:254
Definition concurrentqueue.h:548
Definition concurrentqueue.h:522
Definition concurrentqueue.h:618
Definition concurrentqueue.h:624
Definition concurrentqueue.h:83
Definition concurrentqueue.h:307