diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/RedisQueueMessageDrivenEndpointTests-context.xml b/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/RedisQueueMessageDrivenEndpointTests-context.xml index c32aac1c79..38b1ec6a63 100644 --- a/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/RedisQueueMessageDrivenEndpointTests-context.xml +++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/RedisQueueMessageDrivenEndpointTests-context.xml @@ -12,13 +12,17 @@ + + - @@ -28,8 +32,9 @@ - diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/RedisQueueMessageDrivenEndpointTests.java b/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/RedisQueueMessageDrivenEndpointTests.java index 883b72cc97..37f0dd6a88 100644 --- a/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/RedisQueueMessageDrivenEndpointTests.java +++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/RedisQueueMessageDrivenEndpointTests.java @@ -30,6 +30,8 @@ import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; +import org.junit.After; +import org.junit.Before; import org.junit.Ignore; import org.junit.Test; import org.junit.runner.RunWith; @@ -41,6 +43,7 @@ import org.springframework.beans.factory.DisposableBean; import org.springframework.beans.factory.InitializingBean; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.ApplicationEvent; +import org.springframework.context.Lifecycle; import org.springframework.data.redis.RedisConnectionFailureException; import org.springframework.data.redis.RedisSystemException; import org.springframework.data.redis.connection.RedisConnectionFactory; @@ -78,24 +81,38 @@ import org.springframework.util.ClassUtils; @DirtiesContext public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests { + public static final String TEST_QUEUE = "testQueue"; + @Autowired private RedisConnectionFactory connectionFactory; @Autowired private PollableChannel fromChannel; + @Autowired + private Lifecycle fromChannelEndpoint; + @Autowired private MessageChannel symmetricalInputChannel; + @Autowired + private Lifecycle symmetricalRedisChannelEndpoint; + @Autowired private PollableChannel symmetricalOutputChannel; + @Before + @After + public void setUpTearDown() { + RedisTemplate redisTemplate = new RedisTemplate<>(); + redisTemplate.setConnectionFactory(this.connectionFactory); + redisTemplate.delete(TEST_QUEUE); + } + @Test @RedisAvailable @SuppressWarnings("unchecked") public void testInt3014Default() { - String queueName = "si.test.redisQueueInboundChannelAdapterTests"; - RedisTemplate redisTemplate = new RedisTemplate<>(); redisTemplate.setConnectionFactory(this.connectionFactory); redisTemplate.setEnableDefaultSerializer(false); @@ -105,16 +122,16 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests { String payload = "testing"; - redisTemplate.boundListOps(queueName).leftPush(payload); + redisTemplate.boundListOps(TEST_QUEUE).leftPush(payload); Date payload2 = new Date(); - redisTemplate.boundListOps(queueName).leftPush(payload2); + redisTemplate.boundListOps(TEST_QUEUE).leftPush(payload2); PollableChannel channel = new QueueChannel(); RedisQueueMessageDrivenEndpoint endpoint = - new RedisQueueMessageDrivenEndpoint(queueName, this.connectionFactory); + new RedisQueueMessageDrivenEndpoint(TEST_QUEUE, this.connectionFactory); endpoint.setBeanFactory(Mockito.mock(BeanFactory.class)); endpoint.setBeanClassLoader(ClassUtils.getDefaultClassLoader()); endpoint.setOutputChannel(channel); @@ -137,8 +154,6 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests { @RedisAvailable @SuppressWarnings("unchecked") public void testInt3014ExpectMessageTrue() { - String queueName = "si.test.redisQueueInboundChannelAdapterTests2"; - RedisTemplate redisTemplate = new RedisTemplate<>(); redisTemplate.setConnectionFactory(this.connectionFactory); redisTemplate.setEnableDefaultSerializer(false); @@ -148,16 +163,16 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests { Message message = MessageBuilder.withPayload("testing").build(); - redisTemplate.boundListOps(queueName).leftPush(message); + redisTemplate.boundListOps(TEST_QUEUE).leftPush(message); - redisTemplate.boundListOps(queueName).leftPush("test"); + redisTemplate.boundListOps(TEST_QUEUE).leftPush("test"); PollableChannel channel = new QueueChannel(); PollableChannel errorChannel = new QueueChannel(); RedisQueueMessageDrivenEndpoint endpoint = - new RedisQueueMessageDrivenEndpoint(queueName, this.connectionFactory); + new RedisQueueMessageDrivenEndpoint(TEST_QUEUE, this.connectionFactory); endpoint.setBeanFactory(Mockito.mock(BeanFactory.class)); endpoint.setExpectMessage(true); endpoint.setOutputChannel(channel); @@ -186,26 +201,29 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests { @Test @RedisAvailable public void testInt3017IntegrationInbound() { + this.fromChannelEndpoint.start(); String payload = new Date().toString(); RedisTemplate redisTemplate = new StringRedisTemplate(); redisTemplate.setConnectionFactory(this.connectionFactory); redisTemplate.afterPropertiesSet(); - redisTemplate.boundListOps("si.test.Int3017IntegrationInbound") + redisTemplate.boundListOps(TEST_QUEUE) .leftPush("{\"payload\":\"" + payload + "\",\"headers\":{}}"); Message receive = this.fromChannel.receive(10000); assertThat(receive).isNotNull(); assertThat(receive.getPayload()).isEqualTo(payload); + this.fromChannelEndpoint.stop(); } @Test @RedisAvailable public void testInt3017IntegrationSymmetrical() { + this.symmetricalRedisChannelEndpoint.start(); UUID payload = UUID.randomUUID(); Message message = MessageBuilder.withPayload(payload) - .setHeader("redis_queue", "si.test.Int3017IntegrationSymmetrical") + .setHeader("redis_queue", TEST_QUEUE) .build(); this.symmetricalInputChannel.send(message); @@ -213,14 +231,13 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests { Message receive = this.symmetricalOutputChannel.receive(10000); assertThat(receive).isNotNull(); assertThat(receive.getPayload()).isEqualTo(payload); + this.symmetricalRedisChannelEndpoint.stop(); } @Test @RedisAvailable @SuppressWarnings("unchecked") public void testInt3442ProperlyStop() throws Exception { - String queueName = "si.test.testInt3442ProperlyStopTest"; - RedisTemplate redisTemplate = new RedisTemplate<>(); redisTemplate.setConnectionFactory(this.connectionFactory); redisTemplate.setEnableDefaultSerializer(false); @@ -228,11 +245,11 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests { redisTemplate.setValueSerializer(new JdkSerializationRedisSerializer()); redisTemplate.afterPropertiesSet(); - while (redisTemplate.boundListOps(queueName).rightPop() != null) { + while (redisTemplate.boundListOps(TEST_QUEUE).rightPop() != null) { // drain } - RedisQueueMessageDrivenEndpoint endpoint = new RedisQueueMessageDrivenEndpoint(queueName, + RedisQueueMessageDrivenEndpoint endpoint = new RedisQueueMessageDrivenEndpoint(TEST_QUEUE, this.connectionFactory); BoundListOperations boundListOperations = TestUtils.getPropertyValue(endpoint, "boundListOperations", BoundListOperations.class); @@ -252,7 +269,7 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests { waitListening(endpoint); dfa.setPropertyValue("listening", false); - redisTemplate.boundListOps(queueName).leftPush("foo"); + redisTemplate.boundListOps(TEST_QUEUE).leftPush("foo"); CountDownLatch stopLatch = new CountDownLatch(1); @@ -271,14 +288,13 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests { @RedisAvailable @Ignore("LettuceConnectionFactory doesn't support proper reinitialization after 'destroy()'") public void testInt3196Recovery() throws Exception { - String queueName = "test.si.Int3196Recovery"; QueueChannel channel = new QueueChannel(); final List exceptionEvents = new ArrayList<>(); final CountDownLatch exceptionsLatch = new CountDownLatch(2); - RedisQueueMessageDrivenEndpoint endpoint = new RedisQueueMessageDrivenEndpoint(queueName, + RedisQueueMessageDrivenEndpoint endpoint = new RedisQueueMessageDrivenEndpoint(TEST_QUEUE, this.connectionFactory); endpoint.setBeanFactory(Mockito.mock(BeanFactory.class)); endpoint.setApplicationEventPublisher(event -> { @@ -315,7 +331,7 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests { String payload = "testing"; - redisTemplate.boundListOps(queueName).leftPush(payload); + redisTemplate.boundListOps(TEST_QUEUE).leftPush(payload); Message receive = channel.receive(10000); assertThat(receive).isNotNull(); @@ -328,8 +344,6 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests { @RedisAvailable @SuppressWarnings("unchecked") public void testInt3932ReadFromLeft() { - String queueName = "si.test.redisQueueInboundChannelAdapterTests3932"; - RedisTemplate redisTemplate = new RedisTemplate<>(); redisTemplate.setConnectionFactory(this.connectionFactory); redisTemplate.setEnableDefaultSerializer(false); @@ -339,16 +353,16 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests { String payload = "testing"; - redisTemplate.boundListOps(queueName).rightPush(payload); + redisTemplate.boundListOps(TEST_QUEUE).rightPush(payload); Date payload2 = new Date(); - redisTemplate.boundListOps(queueName).rightPush(payload2); + redisTemplate.boundListOps(TEST_QUEUE).rightPush(payload2); PollableChannel channel = new QueueChannel(); RedisQueueMessageDrivenEndpoint endpoint = - new RedisQueueMessageDrivenEndpoint(queueName, this.connectionFactory); + new RedisQueueMessageDrivenEndpoint(TEST_QUEUE, this.connectionFactory); endpoint.setBeanFactory(Mockito.mock(BeanFactory.class)); endpoint.setBeanClassLoader(ClassUtils.getDefaultClassLoader()); endpoint.setOutputChannel(channel);