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
This commit is contained in:
Aleksey Krichevskiy
2022-08-16 19:17:23 +03:00
committed by Artem Bilan
parent 4de6c0f0c9
commit 217c675caa
6 changed files with 38 additions and 20 deletions

View File

@@ -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'

View File

@@ -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");
}
}

View File

@@ -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;
}

View File

@@ -777,6 +777,14 @@
<xsd:union memberTypes="messageDeletionPolicy xsd:string"/>
</xsd:simpleType>
</xsd:attribute>
<xsd:attribute name="fail-on-missing-queue" type="xsd:boolean">
<xsd:annotation>
<xsd:documentation>
Configures that application should fail on startup if declared queue does not exist.
Default is to ignore missing queues.
</xsd:documentation>
</xsd:annotation>
</xsd:attribute>
</xsd:extension>
</xsd:complexContent>
</xsd:complexType>

View File

@@ -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">
<aws-messaging:sqs-async-client id="sqs"/>
<bean id="resourceIdResolver" class="org.mockito.Mockito" factory-method="mock">
<constructor-arg value="io.awspring.cloud.core.env.ResourceIdResolver"/>
</bean>
@@ -17,6 +15,10 @@
<constructor-arg value="org.springframework.core.task.AsyncTaskExecutor"/>
</bean>
<bean id="destinationResolver" class="org.mockito.Mockito" factory-method="mock">
<constructor-arg value="org.springframework.messaging.core.DestinationResolver"/>
</bean>
<bean class="org.springframework.integration.aws.config.xml.SqsMessageDrivenChannelAdapterParserTests"/>
<int-aws:sqs-message-driven-channel-adapter sqs="sqs"
@@ -34,6 +36,7 @@
send-timeout="2000"
queue-stop-timeout="11000"
destination-resolver="destinationResolver"
resource-id-resolver="resourceIdResolver"/>
resource-id-resolver="resourceIdResolver"
fail-on-missing-queue="true"/>
</beans>

View File

@@ -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);
}
}