diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/core/RabbitTemplate.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/core/RabbitTemplate.java index d06ca27d..ebedd8fd 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/core/RabbitTemplate.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/core/RabbitTemplate.java @@ -33,6 +33,7 @@ import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; +import org.springframework.amqp.AmqpConnectException; import org.springframework.amqp.AmqpException; import org.springframework.amqp.AmqpIllegalStateException; import org.springframework.amqp.AmqpRejectAndDontRequeueException; @@ -665,15 +666,17 @@ public class RabbitTemplate extends RabbitAccessor implements BeanFactoryAware, } if (this.replyAddress == null || Address.AMQ_RABBITMQ_REPLY_TO.equals(this.replyAddress)) { try { - execute(new ChannelCallback() { + return execute(new ChannelCallback() { @Override - public Void doInRabbit(Channel channel) throws Exception { + public Boolean doInRabbit(Channel channel) throws Exception { channel.queueDeclarePassive(Address.AMQ_RABBITMQ_REPLY_TO); - return null; + return true; } }); - return true; + } + catch (AmqpConnectException ex) { + throw ex; } catch (Exception e) { if (this.replyAddress != null) { diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplateTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplateTests.java index 2e08bd1f..ace4b1d0 100644 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplateTests.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplateTests.java @@ -17,10 +17,13 @@ package org.springframework.amqp.rabbit.core; import static org.hamcrest.Matchers.containsString; +import static org.hamcrest.Matchers.instanceOf; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertSame; import static org.junit.Assert.assertThat; import static org.junit.Assert.assertTrue; +import static org.mockito.BDDMockito.willThrow; import static org.mockito.Matchers.any; import static org.mockito.Matchers.anyString; import static org.mockito.Mockito.doAnswer; @@ -44,6 +47,7 @@ import org.mockito.invocation.InvocationOnMock; import org.mockito.stubbing.Answer; import org.springframework.amqp.AmqpAuthenticationException; +import org.springframework.amqp.AmqpConnectException; import org.springframework.amqp.core.Message; import org.springframework.amqp.core.MessageProperties; import org.springframework.amqp.core.ReceiveAndReplyCallback; @@ -123,6 +127,7 @@ public class RabbitTemplateTests { template.setChannelTransacted(true); txTemplate.execute(new TransactionCallback() { + @Override public Object doInTransaction(TransactionStatus status) { template.convertAndSend("foo", "bar"); @@ -130,6 +135,7 @@ public class RabbitTemplateTests { } }); txTemplate.execute(new TransactionCallback() { + @Override public Object doInTransaction(TransactionStatus status) { template.convertAndSend("baz", "qux"); @@ -188,6 +194,7 @@ public class RabbitTemplateTests { final AtomicReference consumer = new AtomicReference(); doAnswer(new Answer() { + @Override public Object answer(InvocationOnMock invocation) throws Throwable { consumer.set((Consumer) invocation.getArguments()[6]); @@ -224,6 +231,23 @@ public class RabbitTemplateTests { assertEquals(3, count.get()); } + @Test + public void testEvaluateDirectReplyToWithConnectException() { + org.springframework.amqp.rabbit.connection.ConnectionFactory mockConnectionFactory = + mock(org.springframework.amqp.rabbit.connection.ConnectionFactory.class); + willThrow(new AmqpConnectException(null)).given(mockConnectionFactory).createConnection(); + RabbitTemplate template = new RabbitTemplate(mockConnectionFactory); + + try { + template.convertSendAndReceive("foo"); + } + catch (Exception ex) { + assertThat(ex, instanceOf(AmqpConnectException.class)); + } + + assertFalse(TestUtils.getPropertyValue(template, "evaluatedFastReplyTo", Boolean.class)); + } + @Test public void testRecovery() throws Exception { ConnectionFactory mockConnectionFactory = mock(ConnectionFactory.class);