diff options
| author | Robert Gemmell <robbie@apache.org> | 2011-06-07 11:18:41 +0000 |
|---|---|---|
| committer | Robert Gemmell <robbie@apache.org> | 2011-06-07 11:18:41 +0000 |
| commit | 67788db66c38665e454ecbf977e1d376bf7eea97 (patch) | |
| tree | fb6ee01c54b5350abc109315ae03b6e254cc22b1 /qpid/java/broker/src/main | |
| parent | 174789ad53876088c38ea365a65ae54ad2c48867 (diff) | |
| download | qpid-python-67788db66c38665e454ecbf977e1d376bf7eea97.tar.gz | |
QPID-3219: update handling of QueueEntries to exclude use of entries in the intermediate 'dequeued' state, simplify logic in general.
Applied patch from Oleksandr Rudyy <orudyy@gmail.com>
git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1132959 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'qpid/java/broker/src/main')
5 files changed, 51 insertions, 16 deletions
diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/AMQPriorityQueue.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/AMQPriorityQueue.java index b6e97e08fb..371ae0de50 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/AMQPriorityQueue.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/AMQPriorityQueue.java @@ -60,7 +60,7 @@ public class AMQPriorityQueue extends SimpleAMQQueue { // check that all subscriptions are not in advance of the entry SubscriptionList.SubscriptionNodeIterator subIter = _subscriptionList.iterator(); - while(subIter.advance() && !entry.isAcquired()) + while(subIter.advance() && entry.isAvailable()) { final Subscription subscription = subIter.getNode().getSubscription(); if(!subscription.isClosed()) @@ -70,7 +70,7 @@ public class AMQPriorityQueue extends SimpleAMQQueue { QueueEntry subnode = context._lastSeenEntry; QueueEntry released = context._releasedEntry; - while(subnode != null && entry.compareTo(subnode) < 0 && !entry.isAcquired() && (released == null || released.compareTo(entry) < 0)) + while(subnode != null && entry.compareTo(subnode) < 0 && entry.isAvailable() && (released == null || released.compareTo(entry) < 0)) { if(QueueContext._releasedUpdater.compareAndSet(context,released,entry)) { diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/QueueEntry.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/QueueEntry.java index 79ede2694e..88349586c3 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/QueueEntry.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/QueueEntry.java @@ -52,6 +52,17 @@ public interface QueueEntry extends Comparable<QueueEntry>, Filterable } public abstract State getState(); + + /** + * Returns true if state is either DEQUEUED or DELETED. + * + * @return true if state is either DEQUEUED or DELETED. + */ + public boolean isDispensed() + { + State currentState = getState(); + return currentState == State.DEQUEUED || currentState == State.DELETED; + } } @@ -207,4 +218,18 @@ public interface QueueEntry extends Comparable<QueueEntry>, Filterable void addStateChangeListener(StateChangeListener listener); boolean removeStateChangeListener(StateChangeListener listener); + + /** + * Returns true if entry is in DEQUEUED state, otherwise returns false. + * + * @return true if entry is in DEQUEUED state, otherwise returns false + */ + boolean isDequeued(); + + /** + * Returns true if entry is either DEQUED or DELETED state. + * + * @return true if entry is either DEQUED or DELETED state + */ + boolean isDispensed(); } diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/QueueEntryImpl.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/QueueEntryImpl.java index 809ba3277e..bc452d2d72 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/QueueEntryImpl.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/QueueEntryImpl.java @@ -499,7 +499,7 @@ public class QueueEntryImpl implements QueueEntry { QueueEntryImpl next = nextNode(); - while(next != null && next.isDeleted()) + while(next != null && next.isDispensed() ) { final QueueEntryImpl newNext = next.nextNode(); @@ -547,4 +547,16 @@ public class QueueEntryImpl implements QueueEntry return _queueEntryList; } + @Override + public boolean isDequeued() + { + return _state == DEQUEUED_STATE; + } + + @Override + public boolean isDispensed() + { + return _state.isDispensed(); + } + } diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/SimpleAMQQueue.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/SimpleAMQQueue.java index b02d03a1ad..274cb6714a 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/SimpleAMQQueue.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/SimpleAMQQueue.java @@ -629,7 +629,7 @@ public class SimpleAMQQueue implements AMQQueue, Subscription.StateListener // this catches the case where we *just* miss an update int loops = 2; - while (!(entry.isAcquired() || entry.isDeleted()) && loops != 0) + while (entry.isAvailable() && loops != 0) { if (nextNode == null) { @@ -648,7 +648,7 @@ public class SimpleAMQQueue implements AMQQueue, Subscription.StateListener } - if (!(entry.isAcquired() || entry.isDeleted())) + if (entry.isAvailable()) { checkSubscriptionsNotAheadOfDelivery(entry); @@ -942,7 +942,7 @@ public class SimpleAMQQueue implements AMQQueue, Subscription.StateListener while (queueListIterator.advance()) { QueueEntry node = queueListIterator.getNode(); - if (node != null && !node.isDeleted()) + if (node != null && !node.isDispensed()) { entryList.add(node); } @@ -1046,7 +1046,7 @@ public class SimpleAMQQueue implements AMQQueue, Subscription.StateListener while (queueListIterator.advance() && !filter.filterComplete()) { QueueEntry node = queueListIterator.getNode(); - if (!node.isDeleted() && filter.accept(node)) + if (!node.isDispensed() && filter.accept(node)) { entryList.add(node); } @@ -1240,7 +1240,6 @@ public class SimpleAMQQueue implements AMQQueue, Subscription.StateListener if ((messageId >= fromMessageId) && (messageId <= toMessageId) - && !node.isDeleted() && node.acquire()) { dequeueEntry(node); @@ -1270,7 +1269,7 @@ public class SimpleAMQQueue implements AMQQueue, Subscription.StateListener while (noDeletes && queueListIterator.advance()) { QueueEntry node = queueListIterator.getNode(); - if (!node.isDeleted() && node.acquire()) + if (node.acquire()) { dequeueEntry(node); noDeletes = false; @@ -1300,7 +1299,7 @@ public class SimpleAMQQueue implements AMQQueue, Subscription.StateListener while (queueListIterator.advance()) { QueueEntry node = queueListIterator.getNode(); - if (!node.isDeleted() && node.acquire()) + if (node.acquire()) { dequeueEntry(node, txn); if(++count == request) @@ -1654,7 +1653,7 @@ public class SimpleAMQQueue implements AMQQueue, Subscription.StateListener QueueEntry node = getNextAvailableEntry(sub); - if (node != null && !(node.isAcquired() || node.isDeleted())) + if (node != null && node.isAvailable()) { if (sub.hasInterest(node)) { @@ -1715,7 +1714,7 @@ public class SimpleAMQQueue implements AMQQueue, Subscription.StateListener QueueEntry node = (releasedNode != null && lastSeen.compareTo(releasedNode)>=0) ? releasedNode : _entries.next(lastSeen); boolean expired = false; - while (node != null && (node.isAcquired() || node.isDeleted() || (expired = node.expired()) || !sub.hasInterest(node))) + while (node != null && (!node.isAvailable() || (expired = node.expired()) || !sub.hasInterest(node))) { if (expired) { @@ -1884,8 +1883,8 @@ public class SimpleAMQQueue implements AMQQueue, Subscription.StateListener while (queueListIterator.advance()) { QueueEntry node = queueListIterator.getNode(); - // Only process nodes that are not currently deleted - if (!node.isDeleted()) + // Only process nodes that are not currently deleted and not dequeued + if (!node.isDispensed()) { // If the node has exired then aquire it if (node.expired() && node.acquire()) diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/SimpleQueueEntryList.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/SimpleQueueEntryList.java index b97c2c55c5..46baab8c85 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/SimpleQueueEntryList.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/SimpleQueueEntryList.java @@ -1,6 +1,5 @@ package org.apache.qpid.server.queue; -import org.apache.qpid.server.message.InboundMessage; import org.apache.qpid.server.message.ServerMessage; import java.util.concurrent.atomic.AtomicReferenceFieldUpdater; @@ -156,7 +155,7 @@ public class SimpleQueueEntryList implements QueueEntryList if(!atTail()) { QueueEntryImpl nextNode = _lastNode.nextNode(); - while(nextNode.isDeleted() && nextNode.nextNode() != null) + while(nextNode.isDispensed() && nextNode.nextNode() != null) { nextNode = nextNode.nextNode(); } |
