diff --git a/build.gradle b/build.gradle
index d3f7c61..4a833b3 100644
--- a/build.gradle
+++ b/build.gradle
@@ -33,7 +33,7 @@ ext {
junitVersion = '5.8.2'
servletApiVersion = '5.0.0'
log4jVersion = '2.17.2'
- springCloudAwsVersion = '2.4.1'
+ springCloudAwsVersion = '2.4.2'
springIntegrationVersion = '6.0.0-M3'
kinesisClientVersion = '1.14.8'
kinesisProducerVersion = '0.14.12'
diff --git a/src/main/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParser.java b/src/main/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParser.java
index 1e437fd..dab4342 100644
--- a/src/main/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParser.java
+++ b/src/main/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParser.java
@@ -80,6 +80,7 @@ public class SqsMessageDrivenChannelAdapterParser extends AbstractSingleBeanDefi
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "visibility-timeout");
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "wait-time-out");
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "queue-stop-timeout");
+ IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "fail-on-missing-queue");
}
}
diff --git a/src/main/java/org/springframework/integration/aws/inbound/SqsMessageDrivenChannelAdapter.java b/src/main/java/org/springframework/integration/aws/inbound/SqsMessageDrivenChannelAdapter.java
index 6a11d43..69bd577 100644
--- a/src/main/java/org/springframework/integration/aws/inbound/SqsMessageDrivenChannelAdapter.java
+++ b/src/main/java/org/springframework/integration/aws/inbound/SqsMessageDrivenChannelAdapter.java
@@ -105,6 +105,10 @@ public class SqsMessageDrivenChannelAdapter extends MessageProducerSupport imple
this.simpleMessageListenerContainerFactory.setDestinationResolver(destinationResolver);
}
+ public void setFailOnMissingQueue(boolean failOnMissingQueue) {
+ this.simpleMessageListenerContainerFactory.setFailOnMissingQueue(failOnMissingQueue);
+ }
+
public void setQueueStopTimeout(long queueStopTimeout) {
this.queueStopTimeout = queueStopTimeout;
}
diff --git a/src/main/resources/org/springframework/integration/aws/config/spring-integration-aws.xsd b/src/main/resources/org/springframework/integration/aws/config/spring-integration-aws.xsd
index 486ce3d..5b84958 100644
--- a/src/main/resources/org/springframework/integration/aws/config/spring-integration-aws.xsd
+++ b/src/main/resources/org/springframework/integration/aws/config/spring-integration-aws.xsd
@@ -777,6 +777,14 @@
+
+
+
+ Configures that application should fail on startup if declared queue does not exist.
+ Default is to ignore missing queues.
+
+
+
diff --git a/src/test/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParserTests-context.xml b/src/test/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParserTests-context.xml
index fbe7536..51815aa 100644
--- a/src/test/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParserTests-context.xml
+++ b/src/test/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParserTests-context.xml
@@ -7,8 +7,6 @@
http://www.springframework.org/schema/integration/aws https://www.springframework.org/schema/integration/aws/spring-integration-aws.xsd
http://www.springframework.org/schema/cloud/aws/messaging https://www.springframework.org/schema/cloud/aws/messaging/spring-cloud-aws-messaging.xsd">
-
-
@@ -17,6 +15,10 @@
+
+
+
+
+ resource-id-resolver="resourceIdResolver"
+ fail-on-missing-queue="true"/>
diff --git a/src/test/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParserTests.java b/src/test/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParserTests.java
index 7f1809a..2a056d9 100644
--- a/src/test/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParserTests.java
+++ b/src/test/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParserTests.java
@@ -17,8 +17,8 @@
package org.springframework.integration.aws.config.xml;
import static org.assertj.core.api.Assertions.assertThat;
-import static org.mockito.ArgumentMatchers.anyString;
-import static org.mockito.BDDMockito.willThrow;
+import static org.mockito.BDDMockito.any;
+import static org.mockito.BDDMockito.given;
import org.junit.jupiter.api.Test;
import org.mockito.Mockito;
@@ -30,12 +30,13 @@ import org.springframework.integration.aws.inbound.SqsMessageDrivenChannelAdapte
import org.springframework.integration.channel.NullChannel;
import org.springframework.integration.test.util.TestUtils;
import org.springframework.messaging.MessageChannel;
-import org.springframework.messaging.core.DestinationResolutionException;
import org.springframework.messaging.core.DestinationResolver;
import org.springframework.test.annotation.DirtiesContext;
import org.springframework.test.context.junit.jupiter.SpringJUnitConfig;
import com.amazonaws.services.sqs.AmazonSQS;
+import com.amazonaws.services.sqs.AmazonSQSAsync;
+import com.amazonaws.services.sqs.model.GetQueueAttributesResult;
import io.awspring.cloud.core.env.ResourceIdResolver;
import io.awspring.cloud.messaging.listener.SimpleMessageListenerContainer;
import io.awspring.cloud.messaging.listener.SqsMessageDeletionPolicy;
@@ -70,10 +71,10 @@ public class SqsMessageDrivenChannelAdapterParserTests {
private SqsMessageDrivenChannelAdapter sqsMessageDrivenChannelAdapter;
@Bean
- DestinationResolver> destinationResolver() {
- DestinationResolver> destinationResolver = Mockito.mock(DestinationResolver.class);
- willThrow(DestinationResolutionException.class).given(destinationResolver).resolveDestination(anyString());
- return destinationResolver;
+ AmazonSQSAsync sqs() {
+ AmazonSQSAsync sqs = Mockito.mock(AmazonSQSAsync.class);
+ given(sqs.getQueueAttributes(any())).willReturn(new GetQueueAttributesResult());
+ return sqs;
}
@Test
@@ -87,11 +88,13 @@ public class SqsMessageDrivenChannelAdapterParserTests {
assertThat(TestUtils.getPropertyValue(listenerContainer, "destinationResolver"))
.isSameAs(this.destinationResolver);
assertThat(listenerContainer.isRunning()).isFalse();
- assertThat(TestUtils.getPropertyValue(listenerContainer, "maxNumberOfMessages")).isEqualTo(5);
- assertThat(TestUtils.getPropertyValue(listenerContainer, "visibilityTimeout")).isEqualTo(200);
- assertThat(TestUtils.getPropertyValue(listenerContainer, "waitTimeOut")).isEqualTo(40);
- assertThat(TestUtils.getPropertyValue(listenerContainer, "queueStopTimeout")).isEqualTo(11000L);
- assertThat(TestUtils.getPropertyValue(listenerContainer, "autoStartup")).isEqualTo(false);
+ assertThat(listenerContainer)
+ .hasFieldOrPropertyWithValue("maxNumberOfMessages", 5)
+ .hasFieldOrPropertyWithValue("visibilityTimeout", 200)
+ .hasFieldOrPropertyWithValue("waitTimeOut", 40)
+ .hasFieldOrPropertyWithValue("queueStopTimeout", 11000L)
+ .hasFieldOrPropertyWithValue("autoStartup", false)
+ .hasFieldOrPropertyWithValue("failOnMissingQueue", true);
assertThat(this.sqsMessageDrivenChannelAdapter.getPhase()).isEqualTo(100);
assertThat(this.sqsMessageDrivenChannelAdapter.isAutoStartup()).isFalse();
@@ -100,10 +103,9 @@ public class SqsMessageDrivenChannelAdapterParserTests {
.isSameAs(this.errorChannel);
assertThat(TestUtils.getPropertyValue(this.sqsMessageDrivenChannelAdapter, "errorChannel"))
.isSameAs(this.nullChannel);
- assertThat(TestUtils.getPropertyValue(this.sqsMessageDrivenChannelAdapter, "messagingTemplate.sendTimeout"))
- .isEqualTo(2000L);
- assertThat(TestUtils.getPropertyValue(this.sqsMessageDrivenChannelAdapter, "messageDeletionPolicy",
- SqsMessageDeletionPolicy.class)).isEqualTo(SqsMessageDeletionPolicy.NEVER);
+ assertThat(this.sqsMessageDrivenChannelAdapter)
+ .hasFieldOrPropertyWithValue("messagingTemplate.sendTimeout", 2000L)
+ .hasFieldOrPropertyWithValue("messageDeletionPolicy", SqsMessageDeletionPolicy.NEVER);
}
}