Fix messages not received if listening on channel and pattern

DATAREDIS-165

Infinite loop in PatternSubscriptionTask causing
broken pipes and unreceived messages
This commit is contained in:
Jennifer Hickey
2013-04-04 11:26:12 -07:00
parent 78ad1a7cd4
commit 2291edae67
2 changed files with 29 additions and 1 deletions

View File

@@ -662,7 +662,7 @@ public class RedisMessageListenerContainer implements InitializingBean, Disposab
// wait for subscription to be initialized
boolean done = false;
// wait 3 rounds for subscription to be initialized
for (int i = 0; i < ROUNDS || done; i++) {
for (int i = 0; i < ROUNDS && !done; i++) {
if (connection != null) {
synchronized (localMonitor) {
if (connection != null && connection.isSubscribed()) {

View File

@@ -22,6 +22,8 @@ import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertTrue;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.HashSet;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Set;
@@ -54,6 +56,8 @@ public class PubSubResubscribeTests {
protected RedisMessageListenerContainer container;
protected RedisConnectionFactory factory;
@SuppressWarnings("rawtypes")
protected RedisTemplate template;
private final BlockingDeque<String> bag = new LinkedBlockingDeque<String>(99);
@@ -191,6 +195,28 @@ public class PubSubResubscribeTests {
assertTrue(set.contains(anotherPayload2));
}
/**
* Validates the behavior of {@link RedisMessageListenerContainer} when it needs to spin up
* a thread executing its PatternSubscriptionTask
* @throws Exception
*/
@Test
public void testInitializeContainerWithMultipleTopicsIncludingPattern() throws Exception {
container.removeMessageListener(adapter);
container.stop();
container.addMessageListener(adapter,
Arrays.asList(new Topic[] { new ChannelTopic(CHANNEL), new PatternTopic("s*") }));
container.start();
template.convertAndSend("somechannel", "HELLO");
template.convertAndSend(CHANNEL, "WORLD");
Set<String> set = new LinkedHashSet<String>();
set.add(bag.poll(1, TimeUnit.SECONDS));
set.add(bag.poll(1, TimeUnit.SECONDS));
assertEquals(new HashSet<String>(Arrays.asList(new String[] { "HELLO", "WORLD" })), set);
}
private class MessageHandler {
private final BlockingDeque<String> bag;
private final String name;
@@ -199,6 +225,8 @@ public class PubSubResubscribeTests {
this.bag = bag;
this.name = name;
}
@SuppressWarnings("unused")
public void handleMessage(String message) {
System.out.println(name + ": " + message);
bag.add(message);