621template <
class TYPE,
class NODE>
629template <
class TYPE,
class NODE>
634 d_queue_p->markReclaim(d_node_p);
639template <
class TYPE,
class NODE>
661 d_queue_p->popComplete(
true);
680 d_queue_p->releaseAllocateLock();
688template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
692 return state >> k_AVAILABLE_SHIFT;
696template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
697void SingleConsumerQueueImpl<TYPE, ATOMIC_OP, MUTEX, CONDITION>
698 ::incrementUntil(AtomicUint *value,
unsigned int bitValue)
700 unsigned int state = ATOMIC_OP::getUintAcquire(value);
701 if (bitValue != (state & 1)) {
702 unsigned int expState;
705 state = ATOMIC_OP::testAndSwapUintAcqRel(value,
708 }
while (state != expState && (bitValue != (state & 1)));
712template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
713void SingleConsumerQueueImpl<TYPE, ATOMIC_OP, MUTEX, CONDITION>
714 ::markReclaim(Node *node)
721 ATOMIC_OP::addInt64AcqRel(&d_capacity, -1);
723 int nodeState = ATOMIC_OP::swapIntAcqRel(&node->d_state, e_RECLAIM);
724 if (e_WRITABLE_AND_BLOCKED == nodeState) {
728 d_readCondition.signal();
732template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
733void SingleConsumerQueueImpl<TYPE, ATOMIC_OP, MUTEX, CONDITION>
734 ::popComplete(
bool destruct)
737 static_cast<Node *
>(ATOMIC_OP::getPtrAcquire(&d_nextRead));
740 nextRead->d_value.object().~TYPE();
743 ATOMIC_OP::setIntRelease(&nextRead->d_state, e_WRITABLE);
745 ATOMIC_OP::setPtrRelease(&d_nextRead,
746 ATOMIC_OP::getPtrAcquire(&nextRead->d_next));
751 if (ATOMIC_OP::getInt64Acquire(&d_capacity) == available(state)) {
755 d_emptyCondition.broadcast();
759template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
760typename SingleConsumerQueueImpl<TYPE, ATOMIC_OP, MUTEX, CONDITION>::Node *
761 SingleConsumerQueueImpl<TYPE, ATOMIC_OP, MUTEX, CONDITION>
764 if (1 == (ATOMIC_OP::getUintAcquire(&d_pushBackDisabled) & 1)) {
773 k_USE_INC - k_AVAILABLE_INC);
775 if (state < 0 || 0 != (state & k_ALLOCATE_MASK)) {
779 state = ATOMIC_OP::addInt64NvAcqRel(&d_state,
780 k_AVAILABLE_INC - k_USE_INC);
786 if (state >= k_AVAILABLE_INC && 0 == (state & k_ALLOCATE_MASK)) {
791 state = ATOMIC_OP::testAndSwapInt64AcqRel(
794 state + k_USE_INC - k_AVAILABLE_INC);
796 else if ( 0 == ( (state >> k_AVAILABLE_SHIFT)
797 + (state & k_USE_MASK))
798 && 0 == (state & k_ALLOCATE_MASK)) {
809 state = ATOMIC_OP::testAndSwapInt64AcqRel(
812 state + k_ALLOCATE_INC);
813 if (expState == state) {
830 SingleConsumerQueueImpl_AllocateLockGuard<
831 SingleConsumerQueueImpl<TYPE,
834 CONDITION> > guard(
this);
836 Node *a =
static_cast<Node *
>(
837 ATOMIC_OP::getPtrAcquire(&d_nextWrite));
839 Node *b =
static_cast<Node *
>(
840 ATOMIC_OP::getPtrAcquire(&a->d_next));
842 Node *nodes =
static_cast<Node *
>(
843 d_allocator.allocate(
sizeof(Node)
844 * k_ALLOCATION_BATCH_SIZE));
845 for (bsl::size_t i = 0;
846 i < k_ALLOCATION_BATCH_SIZE - 1;
849 ATOMIC_OP::initInt(&n->d_state, e_WRITABLE);
850 ATOMIC_OP::initPointer(&n->d_next, n + 1);
853 Node *n = nodes + k_ALLOCATION_BATCH_SIZE - 1;
854 ATOMIC_OP::initInt(&n->d_state, e_WRITABLE);
855 ATOMIC_OP::initPointer(&n->d_next, b);
858 ATOMIC_OP::setPtrRelease(&a->d_next, nodes);
859 ATOMIC_OP::setPtrRelease(&d_nextWrite, nodes);
861 ATOMIC_OP::addInt64AcqRel(&d_capacity,
862 k_ALLOCATION_BATCH_SIZE);
867 ATOMIC_OP::addInt64AcqRel(
869 k_AVAILABLE_INC * (k_ALLOCATION_BATCH_SIZE - 1));
879 state = ATOMIC_OP::getInt64Acquire(&d_state);
884 }
while (state != expState);
888 static_cast<Node *
>(ATOMIC_OP::getPtrAcquire(&d_nextWrite));
892 expNextWrite = nextWrite;
893 Node *next =
static_cast<Node *
>(ATOMIC_OP::getPtrAcquire(
894 &nextWrite->d_next));
896 nextWrite =
static_cast<Node *
>(ATOMIC_OP::testAndSwapPtrAcqRel(
900 }
while (nextWrite != expNextWrite);
902 ATOMIC_OP::addInt64AcqRel(&d_state, -k_USE_INC);
907template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
909void SingleConsumerQueueImpl<TYPE, ATOMIC_OP, MUTEX, CONDITION>
910 ::releaseAllocateLock()
912 ATOMIC_OP::addInt64AcqRel(&d_state, -k_ALLOCATE_INC);
916template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
924, d_allocator(basicAllocator)
926 ATOMIC_OP::initInt64(&d_capacity, 0);
927 ATOMIC_OP::initInt64(&d_state, 0);
929 ATOMIC_OP::initUint(&d_popFrontDisabled, 0);
930 ATOMIC_OP::initUint(&d_pushBackDisabled, 0);
932 ATOMIC_OP::initPointer(&d_nextWrite, 0);
934 Node *n =
static_cast<Node *
>(d_allocator.
allocate(
sizeof(Node)));
935 ATOMIC_OP::initInt(&n->d_state, e_WRITABLE);
936 ATOMIC_OP::initPointer(&n->d_next, n);
938 ATOMIC_OP::setPtrRelease(&d_nextWrite, n);
939 ATOMIC_OP::setPtrRelease(&d_nextRead, n);
942template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
951, d_allocator(basicAllocator)
953 ATOMIC_OP::initInt64(&d_capacity, 0);
954 ATOMIC_OP::initInt64(&d_state, 0);
956 ATOMIC_OP::initUint(&d_popFrontDisabled, 0);
957 ATOMIC_OP::initUint(&d_pushBackDisabled, 0);
959 ATOMIC_OP::initPointer(&d_nextWrite, 0);
961 Node *nodes =
static_cast<Node *
>(d_allocator.
allocate(
sizeof(Node)
963 for (bsl::size_t i = 0; i < capacity; ++i) {
965 ATOMIC_OP::initInt(&n->d_state, e_WRITABLE);
966 ATOMIC_OP::initPointer(&n->d_next, n + 1);
969 Node *n = nodes + capacity;
970 ATOMIC_OP::initInt(&n->d_state, e_WRITABLE);
971 ATOMIC_OP::initPointer(&n->d_next, nodes);
974 ATOMIC_OP::setPtrRelease(&d_nextWrite, nodes);
975 ATOMIC_OP::setPtrRelease(&d_nextRead, nodes);
977 ATOMIC_OP::addInt64AcqRel(&d_capacity, capacity);
978 ATOMIC_OP::addInt64AcqRel(&d_state, k_AVAILABLE_INC * capacity);
981template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
985 Node *end =
static_cast<Node *
>(ATOMIC_OP::getPtrAcquire(&d_nextWrite));
988 Node *at =
static_cast<Node *
>(ATOMIC_OP::getPtrAcquire(&end->d_next));
992 static_cast<Node *
>(ATOMIC_OP::getPtrAcquire(&at->d_next));
994 if (e_READABLE == ATOMIC_OP::getIntAcquire(&at->d_state)) {
995 at->d_value.object().~TYPE();
1001 if (e_READABLE == ATOMIC_OP::getIntAcquire(&at->d_state)) {
1002 at->d_value.object().~TYPE();
1008template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1012 unsigned int generation = ATOMIC_OP::getUintAcquire(&d_popFrontDisabled);
1013 if (1 == (generation & 1)) {
1018 static_cast<Node *
>(ATOMIC_OP::getPtrAcquire(&d_nextRead));
1019 int nodeState = ATOMIC_OP::getIntAcquire(&nextRead->d_state);
1027 if (e_WRITABLE == nodeState) {
1029 nodeState = ATOMIC_OP::getIntAcquire(&nextRead->d_state);
1030 if (e_WRITABLE == nodeState) {
1032 nodeState = ATOMIC_OP::swapIntAcqRel(&nextRead->d_state,
1033 e_WRITABLE_AND_BLOCKED);
1034 while (e_READABLE != nodeState && e_RECLAIM != nodeState) {
1036 ATOMIC_OP::getUintAcquire(&d_popFrontDisabled)) {
1037 ATOMIC_OP::testAndSwapIntAcqRel(&nextRead->d_state,
1038 e_WRITABLE_AND_BLOCKED,
1042 d_readCondition.wait(&d_readMutex);
1043 nodeState = ATOMIC_OP::getIntAcquire(&nextRead->d_state);
1047 if (e_RECLAIM == nodeState) {
1048 ATOMIC_OP::addInt64AcqRel(&d_capacity, 1);
1051 static_cast<Node *
>(ATOMIC_OP::getPtrAcquire(&d_nextRead));
1052 nodeState = ATOMIC_OP::getIntAcquire(&nextRead->d_state);
1054 }
while (e_RECLAIM == nodeState);
1060 CONDITION> > guard(
this);
1062#if defined(BSLMF_MOVABLEREF_USES_RVALUE_REFERENCES)
1065 *value = nextRead->d_value.object();
1071template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1075 Node *target = pushBackHelper();
1086 Node> proctor(
this, target);
1094 int nodeState = ATOMIC_OP::swapIntAcqRel(&target->d_state, e_READABLE);
1095 if (e_WRITABLE_AND_BLOCKED == nodeState) {
1099 d_readCondition.signal();
1105template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1109 Node *target = pushBackHelper();
1120 Node> proctor(
this, target);
1122 TYPE& dummy = value;
1129 int nodeState = ATOMIC_OP::swapIntAcqRel(&target->d_state, e_READABLE);
1130 if (e_WRITABLE_AND_BLOCKED == nodeState) {
1134 d_readCondition.signal();
1140template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1148 static_cast<Node *
>(ATOMIC_OP::getPtrAcquire(&d_nextRead));
1149 int nodeState = ATOMIC_OP::getIntAcquire(&nextRead->d_state);
1150 while (e_READABLE == nodeState || e_RECLAIM == nodeState) {
1151 if (e_READABLE == nodeState) {
1152 nextRead->d_value.object().~TYPE();
1157 ATOMIC_OP::setIntRelease(&nextRead->d_state, e_WRITABLE);
1159 static_cast<Node *
>(ATOMIC_OP::getPtrAcquire(&nextRead->d_next));
1160 ATOMIC_OP::setPtrRelease(&d_nextRead, nextRead);
1161 nodeState = ATOMIC_OP::getIntAcquire(&nextRead->d_state);
1165 ATOMIC_OP::addInt64AcqRel(&d_capacity, reclaim);
1166 ATOMIC_OP::addInt64AcqRel(&d_state, k_AVAILABLE_INC * count);
1171 d_emptyCondition.broadcast();
1174template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1178 unsigned int generation = ATOMIC_OP::getUintAcquire(&d_popFrontDisabled);
1179 if (1 == (generation & 1)) {
1184 static_cast<Node *
>(ATOMIC_OP::getPtrAcquire(&d_nextRead));
1185 int nodeState = ATOMIC_OP::getIntAcquire(&nextRead->d_state);
1187 while (e_RECLAIM == nodeState) {
1188 ATOMIC_OP::addInt64AcqRel(&d_capacity, 1);
1190 nextRead =
static_cast<Node *
>(ATOMIC_OP::getPtrAcquire(&d_nextRead));
1191 nodeState = ATOMIC_OP::getIntAcquire(&nextRead->d_state);
1194 if (e_READABLE != nodeState) {
1202 CONDITION> > guard(
this);
1204#if defined(BSLMF_MOVABLEREF_USES_RVALUE_REFERENCES)
1207 *value = nextRead->d_value.object();
1213template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1217 return pushBack(value);
1220template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1229template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1233 incrementUntil(&d_popFrontDisabled, 1);
1238 d_readCondition.signal();
1243 d_emptyCondition.broadcast();
1246template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1250 incrementUntil(&d_pushBackDisabled, 1);
1253template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1257 incrementUntil(&d_popFrontDisabled, 0);
1260template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1264 incrementUntil(&d_pushBackDisabled, 0);
1268template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1272 return ATOMIC_OP::getInt64Acquire(&d_capacity) ==
1273 available(ATOMIC_OP::getInt64Acquire(&d_state));
1276template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1282template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1286 return 1 == (ATOMIC_OP::getUintAcquire(&d_popFrontDisabled) & 1);
1289template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1293 return 1 == (ATOMIC_OP::getUintAcquire(&d_pushBackDisabled) & 1);
1296template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1301 return static_cast<bsl::size_t
>(
1303 ? ATOMIC_OP::getInt64Acquire(&d_capacity) - avail
1304 : ATOMIC_OP::getInt64Acquire(&d_capacity));
1307template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1311 unsigned int generation = ATOMIC_OP::getUintAcquire(&d_popFrontDisabled);
1312 if (1 == (generation & 1)) {
1319 while (ATOMIC_OP::getInt64Acquire(&d_capacity) != available(state)) {
1320 if (generation != ATOMIC_OP::getUintAcquire(&d_popFrontDisabled)) {
1323 d_emptyCondition.wait(&d_emptyMutex);
1324 state = ATOMIC_OP::getInt64Acquire(&d_state);
1332template <
class TYPE,
class ATOMIC_OP,
class MUTEX,
class CONDITION>
1336 return d_allocator.allocator();