summaryrefslogtreecommitdiff
path: root/qpid/java/broker/src/main
diff options
context:
space:
mode:
authorRobert Gemmell <robbie@apache.org>2012-05-29 11:37:11 +0000
committerRobert Gemmell <robbie@apache.org>2012-05-29 11:37:11 +0000
commit5a888989bd28402e3271b05bcc32a7410f8d17a2 (patch)
tree469fe7c0a7a484cf2663a70d289e342497303f0b /qpid/java/broker/src/main
parent739f8964cdb5baf786759fb900928206a404d3fa (diff)
downloadqpid-python-5a888989bd28402e3271b05bcc32a7410f8d17a2.tar.gz
QPID-3986: Improved tests and resolved some potential thread-safety issues
Applied patch from Oleksandr Rudyy <orudyy@gmail.com>, Philip Harvey <phil@philharveyonline.com> git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1343675 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'qpid/java/broker/src/main')
-rw-r--r--qpid/java/broker/src/main/java/org/apache/qpid/server/protocol/AMQProtocolEngine.java17
1 files changed, 10 insertions, 7 deletions
diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/protocol/AMQProtocolEngine.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/protocol/AMQProtocolEngine.java
index e12c6fa271..849aa05099 100644
--- a/qpid/java/broker/src/main/java/org/apache/qpid/server/protocol/AMQProtocolEngine.java
+++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/protocol/AMQProtocolEngine.java
@@ -591,7 +591,10 @@ public class AMQProtocolEngine implements ServerProtocolEngine, Managable, AMQPr
public List<AMQChannel> getChannels()
{
- return new ArrayList<AMQChannel>(_channelMap.values());
+ synchronized (_channelMap)
+ {
+ return new ArrayList<AMQChannel>(_channelMap.values());
+ }
}
public AMQChannel getAndAssertChannel(int channelId) throws AMQException
@@ -651,6 +654,11 @@ public class AMQProtocolEngine implements ServerProtocolEngine, Managable, AMQPr
synchronized (_channelMap)
{
_channelMap.put(channel.getChannelId(), channel);
+
+ if(_blocking)
+ {
+ channel.block();
+ }
}
}
@@ -659,11 +667,6 @@ public class AMQProtocolEngine implements ServerProtocolEngine, Managable, AMQPr
_cachedChannels[channelId] = channel;
}
- if(_blocking)
- {
- channel.block();
- }
-
checkForNotification();
}
@@ -790,7 +793,7 @@ public class AMQProtocolEngine implements ServerProtocolEngine, Managable, AMQPr
*/
private void closeAllChannels() throws AMQException
{
- for (AMQChannel channel : _channelMap.values())
+ for (AMQChannel channel : getChannels())
{
channel.close();
}