summaryrefslogtreecommitdiff
path: root/java
diff options
context:
space:
mode:
authorMartin Ritchie <ritchiem@apache.org>2007-04-18 14:37:30 +0000
committerMartin Ritchie <ritchiem@apache.org>2007-04-18 14:37:30 +0000
commit70c4ba4e9c6e9d375d21f4dc64eab0c4ad495a63 (patch)
tree9d24c83d4fc9626b41d76e35d8630725f3a8388a /java
parent5ca1ab359a4e6c30286184752e2e26d693ce3753 (diff)
downloadqpid-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.java49
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)