9#ifndef INCLUDED_BDLCC_SINGLEPRODUCERQUEUEIMPL
10#define INCLUDED_BDLCC_SINGLEPRODUCERQUEUEIMPL
116#include <bdlscm_version.h>
134#include <bsl_cstddef.h>
188template <
class TYPE,
class NODE>
231template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
247 static const int k_POP_YIELD_SPIN = 10;
258 static const int k_AVAILABLE_SHIFT = 24;
261 typedef typename ATOMIC_OP::AtomicTypes::Int AtomicInt;
262 typedef typename ATOMIC_OP::AtomicTypes::Uint AtomicUint;
263 typedef typename ATOMIC_OP::AtomicTypes::Int64 AtomicInt64;
264 typedef typename ATOMIC_OP::AtomicTypes::Pointer AtomicPointer;
271 AtomicPointer d_next;
274 typedef QueueNode<TYPE> Node;
277 AtomicPointer d_nextWrite;
279 AtomicPointer d_nextRead;
285 CONDITION d_readCondition;
288 mutable MUTEX d_emptyMutex;
291 mutable CONDITION d_emptyCondition;
298 AtomicUint d_popFrontDisabled;
302 AtomicUint d_pushBackDisabled;
330 static bool allElementsReserved(bsls::Types::Int64 state);
334 static bool canSupplyBlockedThread(bsls::Types::Int64 state);
339 static bool canSupplyOneBlockedThread(bsls::Types::Int64 state);
346 static bool
isEmpty(bsls::Types::Int64 state);
350 static bool willHaveBlockedThread(bsls::Types::Int64 state);
360 void incrementUntil(AtomicUint *value, unsigned int bitValue);
367 void popComplete(Node *node, bool isEmpty);
373 void popFrontRaw(TYPE* value, bool isEmpty);
378 void releaseAllRaw();
388 bslma::UsesBslmaAllocator);
559 d_queue_p->releaseAllRaw();
575template <
class TYPE,
class NODE>
586template <
class TYPE,
class NODE>
590 d_queue_p->popComplete(d_node_p, d_isEmpty);
598template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
603 return (state >> k_AVAILABLE_SHIFT) <= (state & k_BLOCKED_MASK);
606template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
608bool SingleProducerQueueImpl<TYPE, ATOMIC_OP, MUTEX, CONDITION>::
611 return k_AVAILABLE_INC <= state && (state & k_BLOCKED_MASK);
614template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
616bool SingleProducerQueueImpl<TYPE, ATOMIC_OP, MUTEX, CONDITION>::
619 return k_AVAILABLE_INC == (state & k_AVAILABLE_MASK)
620 && (state & k_BLOCKED_MASK);
623template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
628 return state >> k_AVAILABLE_SHIFT;
631template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
636 return k_AVAILABLE_INC > state;
639template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
641bool SingleProducerQueueImpl<TYPE, ATOMIC_OP, MUTEX, CONDITION>::
644 return (state >> k_AVAILABLE_SHIFT) < (state & k_BLOCKED_MASK);
648template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
649void SingleProducerQueueImpl<TYPE, ATOMIC_OP, MUTEX, CONDITION>
650 ::incrementUntil(AtomicUint *value,
unsigned int bitValue)
652 unsigned int state = ATOMIC_OP::getUintAcquire(value);
653 if (bitValue != (state & 1)) {
654 unsigned int expState;
657 state = ATOMIC_OP::testAndSwapUintAcqRel(value,
660 }
while (state != expState && (bitValue != (state & 1)));
664template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
665void SingleProducerQueueImpl<TYPE, ATOMIC_OP, MUTEX, CONDITION>::
666 popComplete(Node *node,
bool isEmpty)
668 node->d_value.object().~TYPE();
670 ATOMIC_OP::setIntRelease(&node->d_state, e_WRITABLE);
676 d_emptyCondition.broadcast();
681template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
682void SingleProducerQueueImpl<TYPE, ATOMIC_OP, MUTEX, CONDITION>::
683 popFrontRaw(TYPE *value,
687 static_cast<Node *
>(ATOMIC_OP::getPtrAcquire(&d_nextRead));
692 static_cast<Node *
>(ATOMIC_OP::getPtrAcquire(&readFrom->d_next));
695 readFrom =
static_cast<Node *
>(ATOMIC_OP::testAndSwapPtrAcqRel(
699 }
while (readFrom != exp);
701 SingleProducerQueueImpl_PopCompleteGuard<
702 SingleProducerQueueImpl <TYPE,
706 Node> guard(
this, readFrom, isEmpty);
708#if defined(BSLMF_MOVABLEREF_USES_RVALUE_REFERENCES)
711 *value = readFrom->d_value.object();
715template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
716void SingleProducerQueueImpl<TYPE, ATOMIC_OP, MUTEX, CONDITION>::
719 Node *
end =
static_cast<Node *
>(ATOMIC_OP::getPtrAcquire(&d_nextWrite));
722 Node *at =
static_cast<Node *
>(ATOMIC_OP::getPtrAcquire(&
end->d_next));
726 static_cast<Node *
>(ATOMIC_OP::getPtrAcquire(&at->d_next));
728 if (e_WRITABLE != ATOMIC_OP::getIntAcquire(&at->d_state)) {
729 at->d_value.object().~TYPE();
732 d_allocator_p->deallocate(at);
737 if (e_WRITABLE != ATOMIC_OP::getIntAcquire(&at->d_state)) {
738 at->d_value.object().~TYPE();
741 d_allocator_p->deallocate(at);
746template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
753, d_allocator_p(
bslma::Default::allocator(basicAllocator))
755 ATOMIC_OP::initInt64(&d_state, 0);
760 ATOMIC_OP::initUint(&d_popFrontDisabled, 0);
761 ATOMIC_OP::initUint(&d_pushBackDisabled, 0);
763 ATOMIC_OP::initPointer(&d_nextWrite, 0);
769 CONDITION> > proctor(
this);
771 Node *n1 =
static_cast<Node *
>(d_allocator_p->
allocate(
sizeof(Node)));
772 ATOMIC_OP::initInt(&n1->d_state, e_WRITABLE);
773 ATOMIC_OP::initPointer(&n1->d_next, n1);
775 ATOMIC_OP::setPtrRelease(&d_nextWrite, n1);
777 ATOMIC_OP::initPointer(&d_nextRead, n1);
779 Node *n2 =
static_cast<Node *
>(d_allocator_p->
allocate(
sizeof(Node)));
780 ATOMIC_OP::initInt(&n2->d_state, e_WRITABLE);
781 ATOMIC_OP::initPointer(&n2->d_next, n1);
782 ATOMIC_OP::setPtrRelease(&n1->d_next, n2);
787template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
795, d_allocator_p(
bslma::Default::allocator(basicAllocator))
797 ATOMIC_OP::initInt64(&d_state, 0);
802 ATOMIC_OP::initUint(&d_popFrontDisabled, 0);
803 ATOMIC_OP::initUint(&d_pushBackDisabled, 0);
805 ATOMIC_OP::initPointer(&d_nextWrite, 0);
811 CONDITION> > proctor(
this);
813 Node *n1 =
static_cast<Node *
>(d_allocator_p->
allocate(
sizeof(Node)));
814 ATOMIC_OP::initInt(&n1->d_state, e_WRITABLE);
815 ATOMIC_OP::initPointer(&n1->d_next, n1);
817 ATOMIC_OP::setPtrRelease(&d_nextWrite, n1);
819 ATOMIC_OP::initPointer(&d_nextRead, n1);
821 Node *n2 =
static_cast<Node *
>(d_allocator_p->
allocate(
sizeof(Node)));
822 ATOMIC_OP::initInt(&n2->d_state, e_WRITABLE);
823 ATOMIC_OP::initPointer(&n2->d_next, n1);
824 ATOMIC_OP::setPtrRelease(&n1->d_next, n2);
826 capacity = (2 <= capacity ? capacity : 2);
828 for (bsl::size_t i = 2; i < capacity; ++i) {
829 Node *n =
static_cast<Node *
>(d_allocator_p->
allocate(
sizeof(Node)));
830 ATOMIC_OP::initInt(&n->d_state, e_WRITABLE);
831 ATOMIC_OP::initPointer(&n->d_next,
832 ATOMIC_OP::getPtrAcquire(&n2->d_next));
834 ATOMIC_OP::setPtrRelease(&n2->d_next, n);
840template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
848template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
852 unsigned int generation = ATOMIC_OP::getUintAcquire(&d_popFrontDisabled);
853 if (1 == (generation & 1)) {
860 if (willHaveBlockedThread(state)) {
862 state = ATOMIC_OP::getInt64Acquire(&d_state);
863 if (willHaveBlockedThread(state)) {
867 state = ATOMIC_OP::addInt64NvAcqRel(
869 k_AVAILABLE_INC + k_BLOCKED_INC);
871 while (isEmpty(state)) {
873 ATOMIC_OP::getUintAcquire(&d_popFrontDisabled)) {
874 ATOMIC_OP::addInt64AcqRel(&d_state, -k_BLOCKED_INC);
877 d_readCondition.wait(&d_readMutex);
878 state = ATOMIC_OP::getInt64Acquire(&d_state);
881 state = ATOMIC_OP::addInt64NvAcqRel(
883 -(k_AVAILABLE_INC + k_BLOCKED_INC));
885 if (canSupplyBlockedThread(state)) {
886 d_readCondition.signal();
891 popFrontRaw(value, isEmpty(state));
896template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
900 if (1 == (ATOMIC_OP::getUintAcquire(&d_pushBackDisabled) & 1)) {
904 Node *nextWrite =
static_cast<Node *
>(
905 ATOMIC_OP::getPtrAcquire(&d_nextWrite));
907 Node *next =
static_cast<Node *
>(
908 ATOMIC_OP::getPtrAcquire(&nextWrite->d_next));
910 if (e_WRITABLE != ATOMIC_OP::getIntAcquire(&next->d_state)) {
911 Node *n =
static_cast<Node *
>(d_allocator_p->allocate(
sizeof(Node)));
913 ATOMIC_OP::initInt(&n->d_state, e_WRITABLE);
914 ATOMIC_OP::initPointer(&n->d_next, next);
916 ATOMIC_OP::setPtrRelease(&nextWrite->d_next, n);
925 ATOMIC_OP::setIntRelease(&nextWrite->d_state, e_READABLE);
926 ATOMIC_OP::setPtrRelease(&d_nextWrite, next);
931 if (canSupplyOneBlockedThread(state)) {
935 d_readCondition.signal();
941template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
945 if (1 == (ATOMIC_OP::getUintAcquire(&d_pushBackDisabled) & 1)) {
949 Node *nextWrite =
static_cast<Node *
>(
950 ATOMIC_OP::getPtrAcquire(&d_nextWrite));
952 Node *next =
static_cast<Node *
>(
953 ATOMIC_OP::getPtrAcquire(&nextWrite->d_next));
955 if (e_WRITABLE != ATOMIC_OP::getIntAcquire(&next->d_state)) {
956 Node *n =
static_cast<Node *
>(d_allocator_p->allocate(
sizeof(Node)));
958 ATOMIC_OP::initInt(&n->d_state, e_WRITABLE);
959 ATOMIC_OP::initPointer(&n->d_next, next);
961 ATOMIC_OP::setPtrRelease(&nextWrite->d_next, n);
971 ATOMIC_OP::setIntRelease(&nextWrite->d_state, e_READABLE);
972 ATOMIC_OP::setPtrRelease(&d_nextWrite, next);
977 if (canSupplyOneBlockedThread(state)) {
981 d_readCondition.signal();
987template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
991 unsigned int generation = ATOMIC_OP::getUintAcquire(&d_popFrontDisabled);
992 if (1 == (generation & 1)) {
1001 while (willHaveBlockedThread(state)) {
1008 state = ATOMIC_OP::testAndSwapInt64AcqRel(&d_state,
1010 state + k_AVAILABLE_INC);
1011 if (expState == state) {
1014 state += k_AVAILABLE_INC;
1015 if (canSupplyBlockedThread(state)) {
1019 d_readCondition.signal();
1025 popFrontRaw(value, isEmpty(state));
1030template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1034 return pushBack(value);
1037template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1044template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1051 if (allElementsReserved(state)) {
1055 state = ATOMIC_OP::testAndSwapInt64AcqRel(
1058 (state & ~k_AVAILABLE_MASK)
1059 | ( (state & k_BLOCKED_MASK)
1060 << k_AVAILABLE_SHIFT));
1061 }
while (state != expState);
1063 state = (state >> k_AVAILABLE_SHIFT) - (state & k_BLOCKED_MASK);
1067 static_cast<Node *
>(ATOMIC_OP::getPtrAcquire(&d_nextRead));
1071 static_cast<Node *
>(ATOMIC_OP::getPtrAcquire(&readFrom->d_next));
1074 readFrom =
static_cast<Node *
>(ATOMIC_OP::testAndSwapPtrAcqRel(
1078 }
while (readFrom != exp);
1080 readFrom->d_value.object().~TYPE();
1082 ATOMIC_OP::setIntRelease(&readFrom->d_state, e_WRITABLE);
1088 d_emptyCondition.broadcast();
1093template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1097 incrementUntil(&d_popFrontDisabled, 1);
1102 d_readCondition.broadcast();
1107 d_emptyCondition.broadcast();
1110template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1114 incrementUntil(&d_pushBackDisabled, 1);
1117template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1121 incrementUntil(&d_popFrontDisabled, 0);
1124template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1128 incrementUntil(&d_pushBackDisabled, 0);
1132template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1137 return isEmpty(state);
1140template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1146template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1150 return 1 == (ATOMIC_OP::getUintAcquire(&d_popFrontDisabled) & 1);
1153template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1157 return 1 == (ATOMIC_OP::getUintAcquire(&d_pushBackDisabled) & 1);
1160template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1167 return avail >= 0 ?
static_cast<bsl::size_t
>(avail) : 0;
1170template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1174 unsigned int generation = ATOMIC_OP::getUintAcquire(&d_popFrontDisabled);
1175 if (1 == (generation & 1)) {
1182 while (!isEmpty(state)) {
1183 if (generation != ATOMIC_OP::getUintAcquire(&d_popFrontDisabled)) {
1186 d_emptyCondition.wait(&d_emptyMutex);
1187 state = ATOMIC_OP::getInt64Acquire(&d_state);
1195template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1199 return d_allocator_p;
#define BSLMF_NESTED_TRAIT_DECLARATION(t_TYPE, t_TRAIT)
Definition bslmf_nestedtraitdeclaration.h:231
Definition bdlcc_singleproducerqueueimpl.h:189
~SingleProducerQueueImpl_PopCompleteGuard()
Definition bdlcc_singleproducerqueueimpl.h:588
Definition bdlcc_singleproducerqueueimpl.h:149
~SingleProducerQueueImpl_ReleaseAllRawProctor()
Definition bdlcc_singleproducerqueueimpl.h:556
void release()
Definition bdlcc_singleproducerqueueimpl.h:565
Definition bdlcc_singleproducerqueueimpl.h:232
void disablePopFront()
Definition bdlcc_singleproducerqueueimpl.h:1095
bool isFull() const
Definition bdlcc_singleproducerqueueimpl.h:1141
SingleProducerQueueImpl(bslma::Allocator *basicAllocator=0)
Definition bdlcc_singleproducerqueueimpl.h:748
void enablePushBack()
Definition bdlcc_singleproducerqueueimpl.h:1126
SingleProducerQueueImpl(bsl::size_t capacity, bslma::Allocator *basicAllocator=0)
Definition bdlcc_singleproducerqueueimpl.h:789
bool isEmpty() const
Definition bdlcc_singleproducerqueueimpl.h:1134
int waitUntilEmpty() const
Definition bdlcc_singleproducerqueueimpl.h:1172
void disablePushBack()
Definition bdlcc_singleproducerqueueimpl.h:1112
int tryPushBack(bslmf::MovableRef< TYPE > value)
Definition bdlcc_singleproducerqueueimpl.h:1038
void enablePopFront()
Definition bdlcc_singleproducerqueueimpl.h:1119
int popFront(TYPE *value)
Definition bdlcc_singleproducerqueueimpl.h:849
int tryPopFront(TYPE *value)
Definition bdlcc_singleproducerqueueimpl.h:988
int tryPushBack(const TYPE &value)
Definition bdlcc_singleproducerqueueimpl.h:1031
TYPE value_type
Definition bdlcc_singleproducerqueueimpl.h:391
bslma::Allocator * allocator() const
Return the allocator used by this object to supply memory.
Definition bdlcc_singleproducerqueueimpl.h:1197
int pushBack(const TYPE &value)
Definition bdlcc_singleproducerqueueimpl.h:897
bsl::size_t numElements() const
Returns the number of elements currently in this queue.
Definition bdlcc_singleproducerqueueimpl.h:1162
bool isPopFrontDisabled() const
Definition bdlcc_singleproducerqueueimpl.h:1148
int pushBack(bslmf::MovableRef< TYPE > value)
Definition bdlcc_singleproducerqueueimpl.h:942
void removeAll()
Definition bdlcc_singleproducerqueueimpl.h:1045
~SingleProducerQueueImpl()
Destroy this object.
Definition bdlcc_singleproducerqueueimpl.h:842
bool isPushBackDisabled() const
Definition bdlcc_singleproducerqueueimpl.h:1155
Definition bslma_allocator.h:545
virtual void * allocate(size_type size)=0
Definition bslmf_movableref.h:752
Definition bslmt_lockguard.h:234
#define BSLS_IDENT(str)
BSLS_IDENT() - insert string into .comment binary segment (if supported)
Definition bsls_ident.h:238
Definition bdlcc_boundedqueue.h:270
T::iterator end(T &container)
Definition bslstl_iterator.h:1621
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
long long Int64
Definition bsls_types.h:134
Definition bsls_objectbuffer.h:277