From 11310ff1ff20f806ec8c16368720a19d88244dc6 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Mon, 1 Jul 2019 14:13:04 -0400 Subject: [PATCH] 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 --- ...ueueMessageDrivenEndpointTests-context.xml | 7 +++- .../RedisQueueMessageDrivenEndpointTests.java | 41 ++++++++++++++++--- 2 files changed, 41 insertions(+), 7 deletions(-) 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..361cdd497b 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 @@ -16,9 +16,10 @@ - @@ -28,8 +29,10 @@ - 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 3b0c9dd1b1..8d9c94f8bf 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 @@ -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 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 redisTemplate = new StringRedisTemplate(); + redisTemplate.setConnectionFactory(this.connectionFactory); + redisTemplate.afterPropertiesSet(); + redisTemplate.delete(queueName); + this.symmetricalRedisChannelEndpoint.start(); UUID payload = UUID.randomUUID(); Message 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 {