diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/BlockingQueueConsumer.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/BlockingQueueConsumer.java index 47a3d4e2..52fadfd7 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/BlockingQueueConsumer.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/BlockingQueueConsumer.java @@ -214,7 +214,7 @@ public class BlockingQueueConsumer { } passiveDeclareTries = 0; } catch (IOException e) { - if (passiveDeclareTries > 0) { + if (passiveDeclareTries > 0 && channel.isOpen()) { if (logger.isWarnEnabled()) { logger.warn("Reconnect failed; retries left=" + (passiveDeclareTries-1), e); try { diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitBindingIntegrationTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitBindingIntegrationTests.java index 0db3be3e..d8ceb6b2 100644 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitBindingIntegrationTests.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitBindingIntegrationTests.java @@ -63,7 +63,7 @@ public class RabbitBindingIntegrationTests { assertEquals("message", result); } finally { - channel.basicCancel(tag); + consumer.getChannel().basicCancel(tag); } return null; @@ -102,7 +102,7 @@ public class RabbitBindingIntegrationTests { assertEquals("message", result); } finally { - channel.basicCancel(tag); + consumer.getChannel().basicCancel(tag); } return null; @@ -239,6 +239,19 @@ public class RabbitBindingIntegrationTests { accessor.getConnectionFactory(), new DefaultMessagePropertiesConverter(), new ActiveObjectCounter(), AcknowledgeMode.AUTO, true, 1, queue.getName()); consumer.start(); + // wait for consumeOk... + int n = 0; + while (n++ < 100) { + if (consumer.getConsumerTag() == null) { + try { + Thread.sleep(100); + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + break; + } + } + } return consumer; } diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplatePublisherCallbacksIntegrationTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplatePublisherCallbacksIntegrationTests.java index bc28c161..c7c3641e 100644 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplatePublisherCallbacksIntegrationTests.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplatePublisherCallbacksIntegrationTests.java @@ -29,6 +29,7 @@ import java.util.List; import java.util.Map; import java.util.Set; import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; @@ -182,7 +183,7 @@ public class RabbitTemplatePublisherCallbacksIntegrationTests { Connection mockConnection = mock(Connection.class); Channel mockChannel = mock(Channel.class); - when(mockConnectionFactory.newConnection()).thenReturn(mockConnection); + when(mockConnectionFactory.newConnection((ExecutorService) null)).thenReturn(mockConnection); when(mockConnection.isOpen()).thenReturn(true); when(mockConnection.createChannel()).thenReturn(new PublisherCallbackChannelImpl(mockChannel)); @@ -209,7 +210,7 @@ public class RabbitTemplatePublisherCallbacksIntegrationTests { Connection mockConnection = mock(Connection.class); Channel mockChannel = mock(Channel.class); - when(mockConnectionFactory.newConnection()).thenReturn(mockConnection); + when(mockConnectionFactory.newConnection((ExecutorService) null)).thenReturn(mockConnection); when(mockConnection.isOpen()).thenReturn(true); PublisherCallbackChannelImpl channel1 = new PublisherCallbackChannelImpl(mockChannel); PublisherCallbackChannelImpl channel2 = new PublisherCallbackChannelImpl(mockChannel); @@ -278,7 +279,7 @@ public class RabbitTemplatePublisherCallbacksIntegrationTests { Connection mockConnection = mock(Connection.class); Channel mockChannel = mock(Channel.class); - when(mockConnectionFactory.newConnection()).thenReturn(mockConnection); + when(mockConnectionFactory.newConnection((ExecutorService) null)).thenReturn(mockConnection); when(mockConnection.isOpen()).thenReturn(true); when(mockConnection.createChannel()).thenReturn(new PublisherCallbackChannelImpl(mockChannel)); @@ -317,7 +318,7 @@ public class RabbitTemplatePublisherCallbacksIntegrationTests { Connection mockConnection = mock(Connection.class); Channel mockChannel = mock(Channel.class); - when(mockConnectionFactory.newConnection()).thenReturn(mockConnection); + when(mockConnectionFactory.newConnection((ExecutorService) null)).thenReturn(mockConnection); when(mockConnection.isOpen()).thenReturn(true); PublisherCallbackChannelImpl callbackChannel = new PublisherCallbackChannelImpl(mockChannel); when(mockConnection.createChannel()).thenReturn(callbackChannel); diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/MessageListenerContainerLifecycleIntegrationTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/MessageListenerContainerLifecycleIntegrationTests.java index 4ff29028..b6e93e61 100755 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/MessageListenerContainerLifecycleIntegrationTests.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/MessageListenerContainerLifecycleIntegrationTests.java @@ -160,7 +160,10 @@ public class MessageListenerContainerLifecycleIntegrationTests { assertEquals(concurrentConsumers, container.getActiveConsumerCount()); container.stop(); - Thread.sleep(1000L); + int n = 0; + while (n++ < 100 && container.getActiveConsumerCount() > 0) { + Thread.sleep(100); + } assertEquals(0, container.getActiveConsumerCount()); if (!transactional) {