diff options
| author | Martin Ritchie <ritchiem@apache.org> | 2007-04-18 14:37:30 +0000 |
|---|---|---|
| committer | Martin Ritchie <ritchiem@apache.org> | 2007-04-18 14:37:30 +0000 |
| commit | 70c4ba4e9c6e9d375d21f4dc64eab0c4ad495a63 (patch) | |
| tree | 9d24c83d4fc9626b41d76e35d8630725f3a8388a /java | |
| parent | 5ca1ab359a4e6c30286184752e2e26d693ce3753 (diff) | |
| download | qpid-python-70c4ba4e9c6e9d375d21f4dc64eab0c4ad495a63.tar.gz | |
QPID-454 Message 'taken' notion is per message. REVERTED as it just wasn't right.. needs to be refactored.
git-svn-id: https://svn.apache.org/repos/asf/incubator/qpid/branches/M2@530037 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'java')
| -rw-r--r-- | java/broker/src/main/java/org/apache/qpid/server/queue/AMQMessage.java | 49 |
1 files changed, 21 insertions, 28 deletions
diff --git a/java/broker/src/main/java/org/apache/qpid/server/queue/AMQMessage.java b/java/broker/src/main/java/org/apache/qpid/server/queue/AMQMessage.java index 23205758c3..b2046efee3 100644 --- a/java/broker/src/main/java/org/apache/qpid/server/queue/AMQMessage.java +++ b/java/broker/src/main/java/org/apache/qpid/server/queue/AMQMessage.java @@ -80,15 +80,17 @@ public class AMQMessage */ private boolean _immediate; + private AtomicBoolean _taken = new AtomicBoolean(false); private TransientMessageData _transientMessageData = new TransientMessageData(); + private Subscription _takenBySubcription; private Set<Subscription> _rejectedBy = null; - private Map<AMQQueue, AtomicBoolean> _takenMap; - private Map<AMQQueue, Subscription> _takenBySubcriptionMap; + private Map<AMQQueue, AtomicBoolean> _takenMap = new HashMap<AMQQueue, AtomicBoolean>(); + private Map<AMQQueue, Subscription> _takenBySubcriptionMap = new HashMap<AMQQueue, Subscription>(); public boolean isTaken(AMQQueue queue) { - return _takenMap.get(queue).get(); + return _taken.get(); } private final int hashcode = System.identityHashCode(this); @@ -204,8 +206,7 @@ public class AMQMessage _immediate = info.isImmediate(); _transientMessageData.setMessagePublishInfo(info); - _takenMap = null; - _takenBySubcriptionMap = null; + _taken = new AtomicBoolean(false); if (_log.isDebugEnabled()) { @@ -323,11 +324,6 @@ public class AMQMessage // enqueuing the messages ensure that if required the destinations are recorded to a // persistent store - int mapSize = _transientMessageData.getDestinationQueues().size(); - - _takenMap = new HashMap<AMQQueue, AtomicBoolean>(mapSize); - _takenBySubcriptionMap = new HashMap<AMQQueue, Subscription>(mapSize); - for (AMQQueue q : _transientMessageData.getDestinationQueues()) { _takenMap.put(q, new AtomicBoolean(false)); @@ -466,17 +462,14 @@ public class AMQMessage public boolean taken(AMQQueue queue, Subscription sub) { - synchronized (queue) + if (_taken.getAndSet(true)) { - if (_takenMap.get(queue).getAndSet(true)) - { - return true; - } - else - { - _takenBySubcriptionMap.put(queue, sub); - return false; - } + return true; + } + else + { + _takenBySubcription = sub; + return false; } } @@ -486,11 +479,8 @@ public class AMQMessage { _log.trace("Releasing Message:" + debugIdentity()); } - synchronized (queue) - { - _takenMap.get(queue).set(false); - _takenBySubcriptionMap.put(queue, null); - } + _taken.set(false); + _takenBySubcription = null; } public boolean checkToken(Object token) @@ -843,13 +833,16 @@ public class AMQMessage public String toString() { - return "Message[" + debugIdentity() + "]: " + _messageId + "; ref count: " + _referenceCount + "; taken for queues: " + - _takenMap.toString() + " by Subs:" + _takenBySubcriptionMap.toString(); + return "Message[" + debugIdentity() + "]: " + _messageId + "; ref count: " + _referenceCount + "; taken : " + + _taken + " by :" + _takenBySubcription; + +// return "Message[" + debugIdentity() + "]: " + _messageId + "; ref count: " + _referenceCount + "; taken for queues: " + +// _takenMap.toString() + " by Subs:" + _takenBySubcriptionMap.toString(); } public Subscription getDeliveredSubscription(AMQQueue queue) { - return _takenBySubcriptionMap.get(queue); + return _takenBySubcription; } public void reject(Subscription subscription) |
