Add start/stop queues management for SQS Adapter

Provide delegates for the SC-AWS `SimpleMessageListenerContainer` `start/stop/isRunning` method for individual queues.

Mark `SqsMessageDrivenChannelAdapter` with `@ManagedResource` and `@IntegrationManagedResource` to make those operation accessible from JMX/ControlBus
Cover the change with test-case
Provide some upgrades
Fix test for NPE from AWS SDK around `lastModified`: https://github.com/aws/aws-sdk-java/issues/822
This commit is contained in:
Artem Bilan
2016-08-23 17:54:16 -04:00
parent 72821bfb8c
commit 8e2fb5e921
4 changed files with 69 additions and 3 deletions

View File

@@ -3,7 +3,7 @@ buildscript {
maven { url 'http://repo.spring.io/plugins-release' }
}
dependencies {
classpath 'io.spring.gradle:spring-io-plugin:0.0.4.RELEASE'
classpath 'io.spring.gradle:spring-io-plugin:0.0.5.RELEASE'
}
}
@@ -56,7 +56,7 @@ ext {
servletApiVersion = '3.1.0'
slf4jVersion = '1.7.21'
springCloudAwsVersion = '1.1.0.RELEASE'
springIntegrationVersion = '4.2.8.RELEASE'
springIntegrationVersion = '4.2.9.RELEASE'
idPrefix = 'aws'

View File

@@ -32,6 +32,9 @@ import org.springframework.cloud.aws.messaging.listener.SqsMessageDeletionPolicy
import org.springframework.core.task.AsyncTaskExecutor;
import org.springframework.integration.aws.support.AwsHeaders;
import org.springframework.integration.endpoint.MessageProducerSupport;
import org.springframework.integration.support.management.IntegrationManagedResource;
import org.springframework.jmx.export.annotation.ManagedOperation;
import org.springframework.jmx.export.annotation.ManagedResource;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageHeaders;
import org.springframework.messaging.core.DestinationResolver;
@@ -51,6 +54,8 @@ import com.amazonaws.services.sqs.AmazonSQSAsync;
* @see SimpleMessageListenerContainer
* @see QueueMessageHandler
*/
@ManagedResource
@IntegrationManagedResource
public class SqsMessageDrivenChannelAdapter extends MessageProducerSupport implements DisposableBean {
private final SimpleMessageListenerContainerFactory simpleMessageListenerContainerFactory =
@@ -140,6 +145,21 @@ public class SqsMessageDrivenChannelAdapter extends MessageProducerSupport imple
this.listenerContainer.stop();
}
@ManagedOperation
public void stop(String logicalQueueName) {
this.listenerContainer.stop(logicalQueueName);
}
@ManagedOperation
public void start(String logicalQueueName) {
this.listenerContainer.start(logicalQueueName);
}
@ManagedOperation
public boolean isRunning(String logicalQueueName) {
return this.listenerContainer.isRunning(logicalQueueName);
}
@Override
public void destroy() throws Exception {
this.listenerContainer.destroy();

View File

@@ -17,6 +17,7 @@
package org.springframework.integration.aws.inbound;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.AssertionsForClassTypes.assertThatThrownBy;
import static org.mockito.BDDMockito.given;
import static org.mockito.BDDMockito.mock;
import static org.mockito.Matchers.any;
@@ -27,12 +28,16 @@ import org.junit.runner.RunWith;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.integration.annotation.ServiceActivator;
import org.springframework.integration.aws.support.AwsHeaders;
import org.springframework.integration.channel.QueueChannel;
import org.springframework.integration.config.EnableIntegration;
import org.springframework.integration.config.ExpressionControlBusFactoryBean;
import org.springframework.integration.core.MessageProducer;
import org.springframework.integration.test.util.TestUtils;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.PollableChannel;
import org.springframework.messaging.support.GenericMessage;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
@@ -58,6 +63,12 @@ public class SqsMessageDrivenChannelAdapterTests {
@Autowired
private SqsMessageDrivenChannelAdapter sqsMessageDrivenChannelAdapter;
@Autowired
private MessageChannel controlBusInput;
@Autowired
private PollableChannel controlBusOutput;
@Test
public void testSqsMessageDrivenChannelAdapter() {
assertThat(TestUtils.getPropertyValue(this.sqsMessageDrivenChannelAdapter,
@@ -71,6 +82,25 @@ public class SqsMessageDrivenChannelAdapterTests {
assertThat(receive).isNotNull();
assertThat((String) receive.getPayload()).isIn("messageContent", "messageContent2");
assertThat(receive.getHeaders().get(AwsHeaders.QUEUE)).isEqualTo("testQueue");
this.controlBusInput.send(new GenericMessage<>("@sqsMessageDrivenChannelAdapter.stop('testQueue')"));
this.controlBusInput.send(new GenericMessage<>("@sqsMessageDrivenChannelAdapter.isRunning('testQueue')"));
receive = this.controlBusOutput.receive(1000);
assertThat(receive).isNotNull();
assertThat((Boolean) receive.getPayload()).isFalse();
this.controlBusInput.send(new GenericMessage<>("@sqsMessageDrivenChannelAdapter.start('testQueue')"));
this.controlBusInput.send(new GenericMessage<>("@sqsMessageDrivenChannelAdapter.isRunning('testQueue')"));
receive = this.controlBusOutput.receive(1000);
assertThat(receive).isNotNull();
assertThat((Boolean) receive.getPayload()).isTrue();
assertThatThrownBy(() ->
this.controlBusInput.send(new GenericMessage<>("@sqsMessageDrivenChannelAdapter.start('foo')")))
.hasCauseExactlyInstanceOf(IllegalArgumentException.class)
.hasMessageContaining("Queue with name 'foo' does not exist");
}
@Configuration
@@ -111,6 +141,19 @@ public class SqsMessageDrivenChannelAdapterTests {
return adapter;
}
@Bean
@ServiceActivator(inputChannel = "controlBusInput")
public ExpressionControlBusFactoryBean controlBus() {
ExpressionControlBusFactoryBean controlBusFactoryBean = new ExpressionControlBusFactoryBean();
controlBusFactoryBean.setOutputChannel(controlBusOutput());
return controlBusFactoryBean;
}
@Bean
public PollableChannel controlBusOutput() {
return new QueueChannel();
}
}
}

View File

@@ -34,6 +34,7 @@ import java.io.IOException;
import java.io.InputStream;
import java.util.Arrays;
import java.util.Collections;
import java.util.Date;
import java.util.HashMap;
import java.util.LinkedList;
import java.util.List;
@@ -317,7 +318,9 @@ public class S3MessageHandlerTests {
AmazonS3 amazonS3 = mock(AmazonS3.class);
given(amazonS3.putObject(any(PutObjectRequest.class))).willReturn(new PutObjectResult());
given(amazonS3.getObjectMetadata(any(GetObjectMetadataRequest.class))).willReturn(new ObjectMetadata());
ObjectMetadata objectMetadata = new ObjectMetadata();
objectMetadata.setLastModified(new Date());
given(amazonS3.getObjectMetadata(any(GetObjectMetadataRequest.class))).willReturn(objectMetadata);
given(amazonS3.copyObject(any(CopyObjectRequest.class))).willReturn(new CopyObjectResult());
ObjectListing objectListing = spy(new ObjectListing());