diff --git a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java index 902e32fe0..76b11a7d5 100644 --- a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java +++ b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java @@ -605,8 +605,7 @@ public class RabbitMessageChannelBinder extends this.template.setUsePublisherConnection(true); } - private final String exchange = deadLetterExchangeName( - properties.getExtension()); + private final String exchange = deadLetterExchangeName(properties.getExtension()); private final String routingKey = properties.getExtension() .getDeadLetterRoutingKey(); @@ -616,6 +615,8 @@ public class RabbitMessageChannelBinder extends private int maxStackTraceLength = -1; + private Boolean dlxPresent; + @Override public void handleMessage(org.springframework.messaging.Message message) throws MessagingException { Message amqpMessage = StaticMessageHeaderAccessor.getSourceData(message); @@ -628,6 +629,9 @@ public class RabbitMessageChannelBinder extends logger.error("No raw message header in " + message); } else { + if (!checkDlx()) { + return; + } Throwable cause = (Throwable) message.getPayload(); if (!shouldRepublish(cause)) { if (logger.isDebugEnabled()) { @@ -692,6 +696,30 @@ public class RabbitMessageChannelBinder extends } } + private boolean checkDlx() { + if (this.dlxPresent == null) { + if (properties.getExtension().isAutoBindDlq()) { + this.dlxPresent = Boolean.TRUE; + } + else { + this.dlxPresent = this.template.execute(channel -> { + String dlx = deadLetterExchangeName(properties.getExtension()); + try { + channel.exchangeDeclarePassive(dlx); + return Boolean.TRUE; + } + catch (IOException e) { + logger.warn("'republishToDlq' is true, but the '" + + dlx + + "' dead letter exchange is not present; disabling 'republishToDlq'"); + return Boolean.FALSE; + } + }); + } + } + return this.dlxPresent; + } + /** * Traverse the cause tree, stopping at AmqpRejectAndDontRequeueException * or ImmediateAcknowledgeAmqpException. diff --git a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RepublishUnitTests.java b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RepublishUnitTests.java new file mode 100644 index 000000000..873348d63 --- /dev/null +++ b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RepublishUnitTests.java @@ -0,0 +1,76 @@ +/* + * Copyright 2020-2020 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.stream.binder.rabbit; + +import java.io.IOException; +import java.util.Collections; + +import com.rabbitmq.client.Channel; +import org.junit.Test; + +import org.springframework.amqp.core.Message; +import org.springframework.amqp.core.MessageProperties; +import org.springframework.amqp.rabbit.connection.Connection; +import org.springframework.amqp.rabbit.connection.ConnectionFactory; +import org.springframework.boot.autoconfigure.amqp.RabbitProperties; +import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; +import org.springframework.cloud.stream.binder.rabbit.properties.RabbitConsumerProperties; +import org.springframework.cloud.stream.provisioning.ConsumerDestination; +import org.springframework.integration.IntegrationMessageHeaderAccessor; +import org.springframework.messaging.MessageHandler; +import org.springframework.messaging.support.ErrorMessage; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.BDDMockito.given; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; + +/** + * @author Gary Russell + * @since 3.0.6 + * + */ +public class RepublishUnitTests { + + @Test + public void testBadRepublishSetting() throws IOException { + ConnectionFactory cf = mock(ConnectionFactory.class); + Connection conn = mock(Connection.class); + given(cf.createConnection()).willReturn(conn); + Channel channel = mock(Channel.class); + given(channel.isOpen()).willReturn(true); + given(channel.exchangeDeclarePassive("DLX")).willThrow(new IOException()); + given(conn.createChannel(false)).willReturn(channel); + RabbitProperties props = new RabbitProperties(); + RabbitMessageChannelBinder binder = new RabbitMessageChannelBinder(cf, props, null); + RabbitConsumerProperties extension = new RabbitConsumerProperties(); + ExtendedConsumerProperties bindingProps = + new ExtendedConsumerProperties(extension); + MessageHandler handler = binder.getErrorMessageHandler(mock(ConsumerDestination.class), "foo", bindingProps); + ErrorMessage message = new ErrorMessage(new RuntimeException("test"), + Collections.singletonMap(IntegrationMessageHeaderAccessor.SOURCE_DATA, + new Message("foo".getBytes(), new MessageProperties()))); + handler.handleMessage(message); + handler.handleMessage(message); + verify(channel, times(1)).exchangeDeclarePassive("DLX"); + verify(channel, never()).basicPublish(any(), any(), eq(false), any(), any()); + } + +}