From 217c675caa84fd8bb4720608239e79b57be6e379 Mon Sep 17 00:00:00 2001 From: Aleksey Krichevskiy Date: Tue, 16 Aug 2022 19:17:23 +0300 Subject: [PATCH] GH-207: Add SQS `fail-on-missing-queue` Fixes https://github.com/spring-projects/spring-integration-aws/issues/207 * Expose `failOnMissingQueue` flag support for `SqsMessageDrivenChannelAdapter` * Upgrade to `spring-cloud-aws-2.4.2` * `fail-on-missing-queue` property support * `SqsMessageDrivenChannelAdapterParserTests` update **Cherry-pick to `main`** # Conflicts: # build.gradle --- build.gradle | 2 +- .../SqsMessageDrivenChannelAdapterParser.java | 1 + .../SqsMessageDrivenChannelAdapter.java | 4 +++ .../aws/config/spring-integration-aws.xsd | 8 +++++ ...rivenChannelAdapterParserTests-context.xml | 9 +++-- ...essageDrivenChannelAdapterParserTests.java | 34 ++++++++++--------- 6 files changed, 38 insertions(+), 20 deletions(-) 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); } }