Fix race condition around Redis key

https://build.spring.io/browse/INT-FATS5IC-922/
https://build.spring.io/browse/INT-MJATS41-1764/

The `RedisQueueMessageDrivenEndpointTests` may be called from different
CI plans against the same Redis instance.

So, clean up keys before and after every unit test

**Cherry-pick until 4.3.x**

# Conflicts:
#	spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/RedisQueueMessageDrivenEndpointTests.java

# Conflicts:
#	spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/RedisQueueMessageDrivenEndpointTests.java
This commit is contained in:
Artem Bilan
2019-07-01 14:13:04 -04:00
parent eb31825bad
commit 11310ff1ff
2 changed files with 41 additions and 7 deletions

View File

@@ -16,9 +16,10 @@
<int:queue/>
</int:channel>
<int-redis:queue-inbound-channel-adapter queue="si.test.Int3017IntegrationInbound"
<int-redis:queue-inbound-channel-adapter id="fromChannelEndpoint" queue="si.test.Int3017IntegrationInbound"
channel="fromChannel"
expect-message="true"
auto-startup="false"
serializer="testSerializer"/>
<bean id="testSerializer" class="org.springframework.integration.redis.util.CustomJsonSerializer"/>
@@ -28,8 +29,10 @@
<int-redis:queue-outbound-channel-adapter queue-expression="headers.redis_queue"/>
</int:chain>
<int-redis:queue-inbound-channel-adapter queue="si.test.Int3017IntegrationSymmetrical"
<int-redis:queue-inbound-channel-adapter id="symmetricalRedisChannelEndpoint"
queue="si.test.Int3017IntegrationSymmetrical"
channel="symmetricalRedisChannel"
auto-startup="false"
serializer=""/>

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2013-2017 the original author or authors.
* Copyright 2013-2019 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
@@ -47,6 +47,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;
@@ -89,9 +90,15 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests {
@Autowired
private PollableChannel fromChannel;
@Autowired
private Lifecycle fromChannelEndpoint;
@Autowired
private MessageChannel symmetricalInputChannel;
@Autowired
private Lifecycle symmetricalRedisChannelEndpoint;
@Autowired
private PollableChannel symmetricalOutputChannel;
@@ -107,6 +114,7 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests {
redisTemplate.setKeySerializer(new StringRedisSerializer());
redisTemplate.setValueSerializer(new JdkSerializationRedisSerializer());
redisTemplate.afterPropertiesSet();
redisTemplate.delete(queueName);
String payload = "testing";
@@ -136,6 +144,7 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests {
assertEquals(payload2, receive.getPayload());
endpoint.stop();
redisTemplate.delete(queueName);
}
@Test
@@ -150,6 +159,7 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests {
redisTemplate.setKeySerializer(new StringRedisSerializer());
redisTemplate.setValueSerializer(new JdkSerializationRedisSerializer());
redisTemplate.afterPropertiesSet();
redisTemplate.delete(queueName);
Message<?> message = MessageBuilder.withPayload("testing").build();
@@ -188,16 +198,21 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests {
Matchers.containsString("java.lang.String cannot be cast to org.springframework.messaging.Message"));
endpoint.stop();
redisTemplate.delete(queueName);
}
@Test
@RedisAvailable
public void testInt3017IntegrationInbound() throws Exception {
String payload = new Date().toString();
String queueName = "si.test.redisQueueInboundChannelAdapterTests2";
RedisTemplate<String, String> redisTemplate = new StringRedisTemplate();
redisTemplate.setConnectionFactory(this.connectionFactory);
redisTemplate.afterPropertiesSet();
redisTemplate.delete(queueName);
this.fromChannelEndpoint.start();
String payload = new Date().toString();
redisTemplate.boundListOps("si.test.Int3017IntegrationInbound")
.leftPush("{\"payload\":\"" + payload + "\",\"headers\":{}}");
@@ -205,14 +220,22 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests {
Message<?> receive = this.fromChannel.receive(10000);
assertNotNull(receive);
assertEquals(payload, receive.getPayload());
this.fromChannelEndpoint.stop();
redisTemplate.delete(queueName);
}
@Test
@RedisAvailable
public void testInt3017IntegrationSymmetrical() throws Exception {
String queueName = "si.test.Int3017IntegrationSymmetrical";
RedisTemplate<String, String> redisTemplate = new StringRedisTemplate();
redisTemplate.setConnectionFactory(this.connectionFactory);
redisTemplate.afterPropertiesSet();
redisTemplate.delete(queueName);
this.symmetricalRedisChannelEndpoint.start();
UUID payload = UUID.randomUUID();
Message<UUID> message = MessageBuilder.withPayload(payload)
.setHeader("redis_queue", "si.test.Int3017IntegrationSymmetrical")
.setHeader("redis_queue", queueName)
.build();
this.symmetricalInputChannel.send(message);
@@ -220,6 +243,8 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests {
Message<?> receive = this.symmetricalOutputChannel.receive(10000);
assertNotNull(receive);
assertEquals(payload, receive.getPayload());
this.symmetricalRedisChannelEndpoint.stop();
redisTemplate.delete(queueName);
}
@Test
@@ -234,6 +259,7 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests {
redisTemplate.setKeySerializer(new StringRedisSerializer());
redisTemplate.setValueSerializer(new JdkSerializationRedisSerializer());
redisTemplate.afterPropertiesSet();
redisTemplate.delete(queueName);
while (redisTemplate.boundListOps(queueName).rightPop() != null) {
// drain
@@ -271,6 +297,7 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests {
assertTrue(stopLatch.await(21, TimeUnit.SECONDS));
verify(boundListOperations, atLeastOnce()).rightPush(any(byte[].class));
redisTemplate.delete(queueName);
}
@@ -285,7 +312,8 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests {
final CountDownLatch exceptionsLatch = new CountDownLatch(2);
RedisQueueMessageDrivenEndpoint endpoint = new RedisQueueMessageDrivenEndpoint(queueName, this.connectionFactory);
RedisQueueMessageDrivenEndpoint endpoint = new RedisQueueMessageDrivenEndpoint(queueName,
this.connectionFactory);
endpoint.setBeanFactory(Mockito.mock(BeanFactory.class));
endpoint.setApplicationEventPublisher(event -> {
exceptionEvents.add((ApplicationEvent) event);
@@ -328,6 +356,7 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests {
assertEquals(payload, receive.getPayload());
endpoint.stop();
redisTemplate.delete(queueName);
}
@Test
@@ -342,6 +371,7 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests {
redisTemplate.setKeySerializer(new StringRedisSerializer());
redisTemplate.setValueSerializer(new JdkSerializationRedisSerializer());
redisTemplate.afterPropertiesSet();
redisTemplate.delete(queueName);
String payload = "testing";
@@ -372,6 +402,7 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests {
assertEquals(payload2, receive.getPayload());
endpoint.stop();
redisTemplate.delete(queueName);
}
private void waitListening(RedisQueueMessageDrivenEndpoint endpoint) throws InterruptedException {