INT-2354 Fix Race Condition in Tests
Some messages could occasionally be lost because the listener container had not yet subscribed to the queue when we started sending messages to it.
This commit is contained in:
committed by
Mark Fisher
parent
f612328a86
commit
4810d3c634
@@ -16,17 +16,23 @@
|
||||
|
||||
package org.springframework.integration.redis.inbound;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
import org.junit.Test;
|
||||
import org.springframework.data.redis.connection.RedisConnection;
|
||||
import org.springframework.data.redis.connection.jedis.JedisConnectionFactory;
|
||||
import org.springframework.data.redis.core.StringRedisTemplate;
|
||||
import org.springframework.data.redis.listener.RedisMessageListenerContainer;
|
||||
import org.springframework.integration.Message;
|
||||
import org.springframework.integration.channel.QueueChannel;
|
||||
import org.springframework.integration.redis.rules.RedisAvailable;
|
||||
import org.springframework.integration.redis.rules.RedisAvailableTests;
|
||||
import org.springframework.integration.test.util.TestUtils;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
import static org.junit.Assert.fail;
|
||||
|
||||
/**
|
||||
* @author Mark Fisher
|
||||
@@ -34,9 +40,17 @@ import static org.junit.Assert.assertTrue;
|
||||
*/
|
||||
public class RedisInboundChannelAdapterTests extends RedisAvailableTests{
|
||||
|
||||
private final Log logger = LogFactory.getLog(this.getClass());
|
||||
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testRedisInboundChannelAdapter() throws Exception {
|
||||
for (int iteration = 0; iteration < 10; iteration ++) {
|
||||
testRedisInboundChannelAdapterGuts(iteration);
|
||||
}
|
||||
}
|
||||
|
||||
private void testRedisInboundChannelAdapterGuts(int iteration) throws Exception {
|
||||
int numToTest = 10;
|
||||
String redisChannelName = "testRedisInboundChannelAdapterChannel";
|
||||
QueueChannel channel = new QueueChannel();
|
||||
@@ -50,18 +64,21 @@ public class RedisInboundChannelAdapterTests extends RedisAvailableTests{
|
||||
adapter.setOutputChannel(channel);
|
||||
adapter.afterPropertiesSet();
|
||||
adapter.start();
|
||||
|
||||
|
||||
RedisMessageListenerContainer container = waitUntilSubscribed(adapter);
|
||||
|
||||
StringRedisTemplate redisTemplate = new StringRedisTemplate(connectionFactory);
|
||||
redisTemplate.afterPropertiesSet();
|
||||
for (int i = 0; i < numToTest; i++) {
|
||||
redisTemplate.convertAndSend(redisChannelName, "test-" + i);
|
||||
String message = "test-" + i + " iteration " + iteration;
|
||||
redisTemplate.convertAndSend(redisChannelName, message);
|
||||
logger.debug("Sent " + message);
|
||||
}
|
||||
int counter = 0;
|
||||
Thread.sleep(2000);
|
||||
for (int i = 0; i < numToTest; i++) {
|
||||
Message<?> message = channel.receive(5000);
|
||||
if (message == null){
|
||||
throw new RuntimeException("Failed to receive message # " + i);
|
||||
throw new RuntimeException("Failed to receive message # " + i + " iteration " + iteration);
|
||||
}
|
||||
assertNotNull(message);
|
||||
assertTrue(message.getPayload().toString().startsWith("test-"));
|
||||
@@ -69,6 +86,35 @@ public class RedisInboundChannelAdapterTests extends RedisAvailableTests{
|
||||
}
|
||||
assertEquals(numToTest, counter);
|
||||
adapter.stop();
|
||||
container.stop();
|
||||
connectionFactory.destroy();
|
||||
}
|
||||
|
||||
/**
|
||||
* Wait until the container has subscribed to the queue and return a
|
||||
* reference to it, so we can stop it at the end of the test.
|
||||
*/
|
||||
protected RedisMessageListenerContainer waitUntilSubscribed(
|
||||
RedisInboundChannelAdapter adapter) throws Exception {
|
||||
RedisMessageListenerContainer container = (RedisMessageListenerContainer) TestUtils
|
||||
.getPropertyValue(adapter, "container");
|
||||
Object subscriptionTask = TestUtils.getPropertyValue(container, "subscriptionTask");
|
||||
RedisConnection connection = (RedisConnection) TestUtils
|
||||
.getPropertyValue(subscriptionTask, "connection");
|
||||
int n = 0;
|
||||
while (true) {
|
||||
if (n++ > 50) {
|
||||
fail("RMLC Failed to Subscribe");
|
||||
}
|
||||
if (connection.isSubscribed()) {
|
||||
logger.debug("Subscribed OK");
|
||||
break;
|
||||
}
|
||||
logger.debug("Waiting...");
|
||||
Thread.sleep(100);
|
||||
}
|
||||
Thread.sleep(100); // Wait a little longer due to race condition in connection.isSubscribed()
|
||||
return container;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user