diff options
| author | Robert Godfrey <rgodfrey@apache.org> | 2013-06-25 10:20:31 +0000 |
|---|---|---|
| committer | Robert Godfrey <rgodfrey@apache.org> | 2013-06-25 10:20:31 +0000 |
| commit | 59b8d464a2a3b36f0985c10c057e14b284e3bc7c (patch) | |
| tree | 113e07b9b6cb40181f74ae3e3fd032ea2815471a /qpid/java/systests/src/main | |
| parent | e280e8fe6d8b5650f3e66e308047d8036ad941f7 (diff) | |
| download | qpid-python-59b8d464a2a3b36f0985c10c057e14b284e3bc7c.tar.gz | |
QPID-4946 : [Java Broker] closing the broker may result in same message being delivered to multipl competing consumers
git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1496401 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'qpid/java/systests/src/main')
| -rw-r--r-- | qpid/java/systests/src/main/java/org/apache/qpid/test/unit/client/connection/BrokerClosesClientConnectionTest.java | 74 |
1 files changed, 73 insertions, 1 deletions
diff --git a/qpid/java/systests/src/main/java/org/apache/qpid/test/unit/client/connection/BrokerClosesClientConnectionTest.java b/qpid/java/systests/src/main/java/org/apache/qpid/test/unit/client/connection/BrokerClosesClientConnectionTest.java index 2cd7520ae4..4a92728d82 100644 --- a/qpid/java/systests/src/main/java/org/apache/qpid/test/unit/client/connection/BrokerClosesClientConnectionTest.java +++ b/qpid/java/systests/src/main/java/org/apache/qpid/test/unit/client/connection/BrokerClosesClientConnectionTest.java @@ -23,6 +23,11 @@ package org.apache.qpid.test.unit.client.connection; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; +import javax.jms.Message; +import javax.jms.MessageConsumer; +import javax.jms.MessageListener; +import javax.jms.MessageProducer; +import javax.naming.NamingException; import org.apache.qpid.AMQConnectionClosedException; import org.apache.qpid.AMQDisconnectedException; import org.apache.qpid.client.AMQConnection; @@ -45,6 +50,7 @@ public class BrokerClosesClientConnectionTest extends QpidBrokerTestCase private Connection _connection; private boolean _isExternalBroker; private final RecordingExceptionListener _recordingExceptionListener = new RecordingExceptionListener(); + private Session _session; @Override protected void setUp() throws Exception @@ -52,7 +58,7 @@ public class BrokerClosesClientConnectionTest extends QpidBrokerTestCase super.setUp(); _connection = getConnection(); - _connection.createSession(false, Session.AUTO_ACKNOWLEDGE); + _session = _connection.createSession(false, Session.AUTO_ACKNOWLEDGE); _connection.setExceptionListener(_recordingExceptionListener); _isExternalBroker = isExternalBroker(); @@ -140,4 +146,70 @@ public class BrokerClosesClientConnectionTest extends QpidBrokerTestCase return _exception; } } + + + private class Listener implements MessageListener + { + int _messageCount; + + @Override + public synchronized void onMessage(Message message) + { + _messageCount++; + } + + public synchronized int getCount() + { + return _messageCount; + } + } + + public void testNoDeliveryAfterBrokerClose() throws JMSException, NamingException, InterruptedException + { + + Listener listener = new Listener(); + + Session session = _connection.createSession(false, Session.CLIENT_ACKNOWLEDGE); + MessageConsumer consumer1 = session.createConsumer(getTestQueue()); + consumer1.setMessageListener(listener); + + MessageProducer producer = _session.createProducer(getTestQueue()); + producer.send(_session.createTextMessage("test message")); + + _connection.start(); + + + synchronized (listener) + { + long currentTime = System.currentTimeMillis(); + long until = currentTime + 2000l; + while(listener.getCount() == 0 && currentTime < until) + { + listener.wait(until - currentTime); + currentTime = System.currentTimeMillis(); + } + } + assertEquals(1, listener.getCount()); + + Connection connection2 = getConnection(); + Session session2 = connection2.createSession(false, Session.CLIENT_ACKNOWLEDGE); + MessageConsumer consumer2 = session2.createConsumer(getTestQueue()); + consumer2.setMessageListener(listener); + connection2.start(); + + + Connection connection3 = getConnection(); + Session session3 = connection3.createSession(false, Session.CLIENT_ACKNOWLEDGE); + MessageConsumer consumer3 = session3.createConsumer(getTestQueue()); + consumer3.setMessageListener(listener); + connection3.start(); + + assertEquals(1, listener.getCount()); + + stopBroker(); + + assertEquals(1, listener.getCount()); + + + } } |
