9#ifndef INCLUDED_BDLCC_SINGLEPRODUCERSINGLECONSUMERBOUNDEDQUEUE
10#define INCLUDED_BDLCC_SINGLEPRODUCERSINGLECONSUMERBOUNDEDQUEUE
224#include <bdlscm_version.h>
257template <
class TYPE,
class NODE>
294#if defined(BSLS_COMPILERFEATURES_SUPPORT_ALIGNAS)
295class alignas(bslmt::Platform::e_CACHE_LINE_SIZE)
303 typedef unsigned int Uint;
305 typedef typename bsls::AtomicOperations::AtomicTypes::Uint AtomicUint;
306 typedef typename bsls::AtomicOperations::AtomicTypes::Uint64 AtomicUint64;
320 e_READABLE_AND_BLOCKED,
324 e_WRITABLE_AND_EMPTY,
327 e_WRITABLE_AND_BLOCKED
334 template <
class DATA>
341 typedef QueueNode<TYPE> Node;
344 AtomicUint64 d_popIndex;
347 Node *d_popElement_p;
352 const bsl::size_t d_popCapacity;
356 AtomicUint d_popDisabledGeneration;
360 mutable AtomicUint d_emptyCount;
363 AtomicUint d_emptyGeneration;
369 -
sizeof(AtomicUint64)
371 -
sizeof(bsl::size_t)
374 -
sizeof(AtomicUint)];
380 AtomicUint64 d_pushIndex;
383 Node *d_pushElement_p;
388 const bsl::size_t d_pushCapacity;
392 AtomicUint d_pushDisabledGeneration;
397 -
sizeof(AtomicUint64)
399 -
sizeof(bsl::size_t)
400 -
sizeof(AtomicUint)];
442 static void incrementUntil(AtomicUint *value, unsigned int bitValue);
452 void popComplete(Node *node, Uint64 index);
467 int popFrontImp(TYPE *value, bool isTry);
478 int pushBackImp(const TYPE& value, bool isTry);
490 int pushBackImp(bslmf::MovableRef<TYPE> value, bool isTry);
496 void pushComplete(Node *node, Uint64 index);
501 const SingleProducerSingleConsumerBoundedQueue&);
503 const SingleProducerSingleConsumerBoundedQueue&);
508 bslma::UsesBslmaAllocator);
690template <
class TYPE,
class NODE>
692SingleProducerSingleConsumerBoundedQueue_PopCompleteGuard<TYPE, NODE>
693 ::SingleProducerSingleConsumerBoundedQueue_PopCompleteGuard(
703template <
class TYPE,
class NODE>
708 d_queue_p->popComplete(d_node_p, d_index);
720 unsigned int state = AtomicOp::getUintAcquire(value);
721 if (bitValue != (state & 1)) {
722 unsigned int expState;
725 state = AtomicOp::testAndSwapUintAcqRel(value,
728 }
while (state != expState && (bitValue != (state & 1)));
735void SingleProducerSingleConsumerBoundedQueue<TYPE>::popComplete(Node *node,
739 if (index == d_popCapacity) {
742 AtomicOp::setUint64Release(&d_popIndex, index);
744 node->d_value.object().~TYPE();
746 Uint nodeState = AtomicOp::swapUintAcqRel(&node->d_state, e_WRITABLE);
747 if (e_READABLE_AND_BLOCKED == nodeState) {
751 d_pushCondition.signal();
757 nodeState = AtomicOp::testAndSwapUintAcqRel(&d_popElement_p[index].d_state,
759 e_WRITABLE_AND_EMPTY);
760 if (e_WRITABLE == nodeState) {
763 AtomicOp::addUintAcqRel(&d_emptyGeneration, 1);
764 if (0 < AtomicOp::getUintAcquire(&d_emptyCount)) {
768 d_emptyCondition.broadcast();
774int SingleProducerSingleConsumerBoundedQueue<TYPE>::popFrontImp(TYPE *value,
777 Uint64 index = AtomicOp::getUint64Acquire(&d_popIndex);
778 const Uint disabledGen =
779 AtomicOp::getUintAcquire(&d_popDisabledGeneration);
781 if (disabledGen & 1) {
785 Node& node = d_popElement_p[index];
787 Uint nodeState = AtomicOp::getUintAcquire(&node.d_state);
795 if (e_WRITABLE_AND_EMPTY == nodeState) {
801 nodeState = AtomicOp::getUintAcquire(&node.d_state);
802 if (e_WRITABLE_AND_EMPTY == nodeState) {
805 nodeState = AtomicOp::testAndSwapUintAcqRel(
808 e_WRITABLE_AND_BLOCKED);
810 while (( e_WRITABLE_AND_EMPTY == nodeState
811 || e_WRITABLE_AND_BLOCKED == nodeState)
813 AtomicOp::getUintAcquire(&d_popDisabledGeneration)) {
814 int rv = d_popCondition.wait(&d_popMutex);
816 AtomicOp::testAndSwapUintAcqRel(&node.d_state,
817 e_WRITABLE_AND_BLOCKED,
818 e_WRITABLE_AND_EMPTY);
821 nodeState = AtomicOp::getUint(&node.d_state);
827 if ( e_WRITABLE_AND_EMPTY == nodeState
828 || e_WRITABLE_AND_BLOCKED == nodeState) {
829 AtomicOp::testAndSwapUintAcqRel(&node.d_state,
830 e_WRITABLE_AND_BLOCKED,
831 e_WRITABLE_AND_EMPTY);
837 SingleProducerSingleConsumerBoundedQueue_PopCompleteGuard<
838 SingleProducerSingleConsumerBoundedQueue<TYPE>, Node>
839 guard(
this, &node, index);
841#if defined(BSLMF_MOVABLEREF_USES_RVALUE_REFERENCES)
844 *value = node.d_value.object();
851int SingleProducerSingleConsumerBoundedQueue<TYPE>::pushBackImp(
855 Uint64 index = AtomicOp::getUint64Acquire(&d_pushIndex);
856 const Uint disabledGen =
857 AtomicOp::getUintAcquire(&d_pushDisabledGeneration);
859 if (disabledGen & 1) {
863 Node& node = d_pushElement_p[index];
865 Uint nodeState = AtomicOp::getUintAcquire(&node.d_state);
873 if (e_READABLE == nodeState) {
879 nodeState = AtomicOp::getUintAcquire(&node.d_state);
880 if (e_READABLE == nodeState) {
883 nodeState = AtomicOp::testAndSwapUintAcqRel(
886 e_READABLE_AND_BLOCKED);
888 while (( e_READABLE == nodeState
889 || e_READABLE_AND_BLOCKED == nodeState)
891 AtomicOp::getUintAcquire(&d_pushDisabledGeneration)) {
892 int rv = d_pushCondition.wait(&d_pushMutex);
894 AtomicOp::testAndSwapUintAcqRel(&node.d_state,
895 e_READABLE_AND_BLOCKED,
899 nodeState = AtomicOp::getUint(&node.d_state);
905 if ( e_READABLE == nodeState
906 || e_READABLE_AND_BLOCKED == nodeState) {
907 AtomicOp::testAndSwapUintAcqRel(&node.d_state,
908 e_READABLE_AND_BLOCKED,
919 pushComplete(&node, index);
925int SingleProducerSingleConsumerBoundedQueue<TYPE>::pushBackImp(
929 Uint64 index = AtomicOp::getUint64Acquire(&d_pushIndex);
930 const Uint disabledGen =
931 AtomicOp::getUintAcquire(&d_pushDisabledGeneration);
933 if (disabledGen & 1) {
937 Node& node = d_pushElement_p[index];
939 Uint nodeState = AtomicOp::getUintAcquire(&node.d_state);
941 if (e_READABLE == nodeState) {
947 nodeState = AtomicOp::getUintAcquire(&node.d_state);
948 if (e_READABLE == nodeState) {
951 nodeState = AtomicOp::testAndSwapUintAcqRel(
954 e_READABLE_AND_BLOCKED);
956 while (( e_READABLE == nodeState
957 || e_READABLE_AND_BLOCKED == nodeState)
959 AtomicOp::getUintAcquire(&d_pushDisabledGeneration)) {
960 int rv = d_pushCondition.wait(&d_pushMutex);
962 AtomicOp::testAndSwapUintAcqRel(&node.d_state,
963 e_READABLE_AND_BLOCKED,
967 nodeState = AtomicOp::getUint(&node.d_state);
970 if ( e_READABLE == nodeState
971 || e_READABLE_AND_BLOCKED == nodeState) {
972 AtomicOp::testAndSwapUintAcqRel(&node.d_state,
973 e_READABLE_AND_BLOCKED,
985 pushComplete(&node, index);
992void SingleProducerSingleConsumerBoundedQueue<TYPE>::pushComplete(
997 Uint nodeState = AtomicOp::swapUintAcqRel(&node->d_state, e_READABLE);
998 if (e_WRITABLE_AND_BLOCKED == nodeState) {
1001 AtomicOp::addUintAcqRel(&d_emptyGeneration, 1);
1005 d_popCondition.signal();
1007 else if (e_WRITABLE_AND_EMPTY == nodeState) {
1010 AtomicOp::addUintAcqRel(&d_emptyGeneration, 1);
1014 if (index == d_pushCapacity) {
1017 AtomicOp::setUint64Release(&d_pushIndex, index);
1021template <
class TYPE>
1026, d_popCapacity(capacity > 0 ? capacity : 1)
1028, d_pushCapacity(capacity > 0 ? capacity : 1)
1036, d_allocator_p(
bslma::Default::allocator(basicAllocator))
1046 d_popElement_p =
static_cast<Node *
>(
1047 d_allocator_p->
allocate(d_popCapacity *
sizeof(Node)));
1049 d_pushElement_p = d_popElement_p;
1052 for (bsl::size_t i = 1; i < d_popCapacity; ++i) {
1057template <
class TYPE>
1061 if (d_popElement_p) {
1063 d_allocator_p->deallocate(d_popElement_p);
1068template <
class TYPE>
1072 return popFrontImp(value,
false);
1075template <
class TYPE>
1079 return pushBackImp(value,
false);
1082template <
class TYPE>
1090template <
class TYPE>
1093 Uint64 index = AtomicOp::getUint64Acquire(&d_popIndex);
1094 Uint nodeState = AtomicOp::getUintAcquire(
1095 &d_popElement_p[index].d_state);
1097 while (e_READABLE == nodeState || e_READABLE_AND_BLOCKED == nodeState) {
1098 d_popElement_p[index].d_value.object().~TYPE();
1100 AtomicOp::swapUintAcqRel(&d_popElement_p[index].d_state, e_WRITABLE);
1103 if (index == d_popCapacity) {
1107 nodeState = AtomicOp::getUintAcquire(&d_popElement_p[index].d_state);
1114 nodeState = AtomicOp::testAndSwapUintAcqRel(&d_popElement_p[index].d_state,
1116 e_WRITABLE_AND_EMPTY);
1118 if (e_WRITABLE == nodeState) {
1121 AtomicOp::addUintAcqRel(&d_emptyGeneration, 1);
1130 AtomicOp::addUintAcqRel(&d_emptyGeneration, 2);
1133 AtomicOp::setUint64Release(&d_popIndex, index);
1138 d_pushCondition.signal();
1140 if (0 < AtomicOp::getUintAcquire(&d_emptyCount)) {
1144 d_emptyCondition.broadcast();
1148template <
class TYPE>
1152 return popFrontImp(value,
true);
1155template <
class TYPE>
1160 return pushBackImp(value,
true);
1163template <
class TYPE>
1173template <
class TYPE>
1177 incrementUntil(&d_popDisabledGeneration, 1);
1182 d_popCondition.broadcast();
1184 if (0 < AtomicOp::getUintAcquire(&d_emptyCount)) {
1188 d_emptyCondition.broadcast();
1192template <
class TYPE>
1196 incrementUntil(&d_pushDisabledGeneration, 1);
1201 d_pushCondition.broadcast();
1204template <
class TYPE>
1208 incrementUntil(&d_popDisabledGeneration, 0);
1211template <
class TYPE>
1215 incrementUntil(&d_pushDisabledGeneration, 0);
1219template <
class TYPE>
1223 return d_popCapacity;
1226template <
class TYPE>
1230 return 0 == (AtomicOp::getUintAcquire(&d_emptyGeneration) & 1);
1233template <
class TYPE>
1237 Node& node = d_pushElement_p[AtomicOp::getUint64Acquire(
1239 Uint nodeState = AtomicOp::getUintAcquire(&node.d_state);
1241 return e_READABLE == nodeState || e_READABLE_AND_BLOCKED == nodeState;
1244template <
class TYPE>
1248 return 1 == (AtomicOp::getUintAcquire(&d_popDisabledGeneration) & 1);
1251template <
class TYPE>
1255 return 1 == (AtomicOp::getUintAcquire(&d_pushDisabledGeneration) & 1);
1258template <
class TYPE>
1262 Uint64 popIndex = AtomicOp::getUint64Acquire(&d_popIndex);
1263 Uint64 pushIndex = AtomicOp::getUint64Acquire(&d_pushIndex);
1264 Node& node = d_pushElement_p[pushIndex];
1265 Uint nodeState = AtomicOp::getUintAcquire(&node.d_state);
1267 if (e_READABLE == nodeState || e_READABLE_AND_BLOCKED == nodeState) {
1268 return d_popCapacity;
1271 return static_cast<bsl::size_t
>( pushIndex >= popIndex
1272 ? pushIndex - popIndex
1273 : pushIndex + d_popCapacity - popIndex);
1276template <
class TYPE>
1279 AtomicOp::addUintAcqRel(&d_emptyCount, 1);
1281 const Uint initEmptyGen = AtomicOp::getUintAcquire(&d_emptyGeneration);
1283 const Uint disabledGen =
1284 AtomicOp::getUintAcquire(&d_popDisabledGeneration);
1286 if (disabledGen & 1) {
1287 AtomicOp::addUintAcqRel(&d_emptyCount, -1);
1291 if (0 == (initEmptyGen & 1)) {
1292 AtomicOp::addUintAcqRel(&d_emptyCount, -1);
1298 Uint emptyGen = AtomicOp::getUintAcquire(&d_emptyGeneration);
1300 while ( initEmptyGen == emptyGen
1302 AtomicOp::getUintAcquire(&d_popDisabledGeneration)) {
1303 int rv = d_emptyCondition.wait(&d_emptyMutex);
1305 AtomicOp::addUintAcqRel(&d_emptyCount, -1);
1308 emptyGen = AtomicOp::getUintAcquire(&d_emptyGeneration);
1311 AtomicOp::addUintAcqRel(&d_emptyCount, -1);
1313 if (initEmptyGen == emptyGen) {
1322template <
class TYPE>
1327 return d_allocator_p;
#define BSLMF_NESTED_TRAIT_DECLARATION(t_TYPE, t_TRAIT)
Definition bslmf_nestedtraitdeclaration.h:231
Definition bdlcc_singleproducersingleconsumerboundedqueue.h:258
~SingleProducerSingleConsumerBoundedQueue_PopCompleteGuard()
Destroy this object and invoke the TYPE::popComplete.
Definition bdlcc_singleproducersingleconsumerboundedqueue.h:706
Definition bdlcc_singleproducersingleconsumerboundedqueue.h:298
void removeAll()
Definition bdlcc_singleproducersingleconsumerboundedqueue.h:1091
int waitUntilEmpty() const
Definition bdlcc_singleproducersingleconsumerboundedqueue.h:1277
bool isPushBackDisabled() const
Definition bdlcc_singleproducersingleconsumerboundedqueue.h:1253
bsl::size_t capacity() const
Definition bdlcc_singleproducersingleconsumerboundedqueue.h:1221
void disablePopFront()
Definition bdlcc_singleproducersingleconsumerboundedqueue.h:1175
int tryPushBack(const TYPE &value)
Definition bdlcc_singleproducersingleconsumerboundedqueue.h:1157
bslma::Allocator * allocator() const
Return the allocator used by this object to supply memory.
Definition bdlcc_singleproducersingleconsumerboundedqueue.h:1324
int pushBack(const TYPE &value)
Definition bdlcc_singleproducersingleconsumerboundedqueue.h:1077
bool isPopFrontDisabled() const
Definition bdlcc_singleproducersingleconsumerboundedqueue.h:1246
TYPE value_type
Definition bdlcc_singleproducersingleconsumerboundedqueue.h:511
bool isEmpty() const
Definition bdlcc_singleproducersingleconsumerboundedqueue.h:1228
~SingleProducerSingleConsumerBoundedQueue()
Destroy this object.
Definition bdlcc_singleproducersingleconsumerboundedqueue.h:1059
void enablePushBack()
Definition bdlcc_singleproducersingleconsumerboundedqueue.h:1213
bool isFull() const
Definition bdlcc_singleproducersingleconsumerboundedqueue.h:1235
void disablePushBack()
Definition bdlcc_singleproducersingleconsumerboundedqueue.h:1194
bsl::size_t numElements() const
Returns the number of elements currently in this queue.
Definition bdlcc_singleproducersingleconsumerboundedqueue.h:1260
int tryPopFront(TYPE *value)
Definition bdlcc_singleproducersingleconsumerboundedqueue.h:1150
int popFront(TYPE *value)
Definition bdlcc_singleproducersingleconsumerboundedqueue.h:1070
void enablePopFront()
Definition bdlcc_singleproducersingleconsumerboundedqueue.h:1206
Definition bslma_allocator.h:545
virtual void * allocate(size_type size)=0
Definition bslmf_movableref.h:752
Definition bslmt_condition.h:220
Definition bslmt_lockguard.h:234
Definition bslmt_mutex.h:317
#define BSLS_IDENT(str)
BSLS_IDENT() - insert string into .comment binary segment (if supported)
Definition bsls_ident.h:238
Definition bdlcc_boundedqueue.h:270
Definition baljsn_encoder_testtypes.h:76
static void moveConstruct(TARGET_TYPE *address, TARGET_TYPE &original, bslma::Allocator *allocator)
Definition bslalg_scalarprimitives.h:1660
static void copyConstruct(TARGET_TYPE *address, const TARGET_TYPE &original, bslma::Allocator *allocator)
Definition bslalg_scalarprimitives.h:1617
static MovableRef< t_TYPE > move(t_TYPE &reference) BSLS_KEYWORD_NOEXCEPT
Definition bslmf_movableref.h:1067
static void yield()
Definition bslmt_threadutil.h:1100
Definition bsls_atomicoperations.h:836
static void initUint64(AtomicTypes::Uint64 *atomicUint, Types::Uint64 initialValue=0)
Definition bsls_atomicoperations.h:2123
static void initUint(AtomicTypes::Uint *atomicUint, unsigned int initialValue=0)
Definition bsls_atomicoperations.h:1924
unsigned long long Uint64
Definition bsls_types.h:139
Definition bsls_objectbuffer.h:277