From 4ce7d1fb45160a2e403a85b1ed5ab7dda26db9b8 Mon Sep 17 00:00:00 2001 From: Costin Leau Date: Fri, 21 Jan 2011 19:34:08 +0200 Subject: [PATCH] + fixed problem caused by message refactoring --- .../connection/jedis/JedisConnection.java | 17 +++++++++------- .../jedis/JedisMessageListener.java | 4 ++-- .../keyvalue/redis/listener/PubSubTests.java | 20 ++++++++++++++++++- 3 files changed, 31 insertions(+), 10 deletions(-) diff --git a/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/jedis/JedisConnection.java b/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/jedis/JedisConnection.java index 7ac70ceb8..9b0057682 100644 --- a/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/jedis/JedisConnection.java +++ b/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/jedis/JedisConnection.java @@ -763,8 +763,9 @@ public class JedisConnection implements RedisConnection { public Boolean getBit(byte[] key, long offset) { try { if (isQueueing()) { - transaction.getbit(key, (int) offset); - return null; + // transaction.getbit(key, (int) offset); + // return null; + throw new UnsupportedOperationException(); } return (jedis.getbit(key, (int) offset) == 0 ? Boolean.FALSE : Boolean.TRUE); } catch (Exception ex) { @@ -776,8 +777,9 @@ public class JedisConnection implements RedisConnection { public void setBit(byte[] key, long offset, boolean value) { try { if (isQueueing()) { - transaction.setbit(key, (int) offset, JedisUtils.asBit(value)); - return; + // transaction.setbit(key, (int) offset, JedisUtils.asBit(value)); + // return; + throw new UnsupportedOperationException(); } jedis.setbit(key, (int) offset, JedisUtils.asBit(value)); } catch (Exception ex) { @@ -873,8 +875,9 @@ public class JedisConnection implements RedisConnection { public Long lInsert(byte[] key, POSITION where, byte[] pivot, byte[] value) { try { if (isQueueing()) { - transaction.linsert(key, JedisUtils.convertPosition(where), pivot, value); - return null; + // transaction.linsert(key, JedisUtils.convertPosition(where), pivot, value); + // return null; + throw new UnsupportedOperationException(); } return jedis.linsert(key, JedisUtils.convertPosition(where), pivot, value); } catch (Exception ex) { @@ -992,7 +995,7 @@ public class JedisConnection implements RedisConnection { if (isQueueing()) { throw new UnsupportedOperationException(); } - return jedis.brpoplpush(srcKey, dstKey, timeout); + return jedis.brpoplpush(srcKey, dstKey, timeout).getBytes(); } catch (Exception ex) { throw convertJedisAccessException(ex); } diff --git a/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/jedis/JedisMessageListener.java b/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/jedis/JedisMessageListener.java index e818e9402..fa0ece097 100644 --- a/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/jedis/JedisMessageListener.java +++ b/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/jedis/JedisMessageListener.java @@ -37,12 +37,12 @@ class JedisMessageListener extends JedisPubSub { @Override public void onMessage(String channel, String message) { - listener.onMessage(new DefaultMessage(message.getBytes(), channel.getBytes()), null); + listener.onMessage(new DefaultMessage(channel.getBytes(), message.getBytes()), null); } @Override public void onPMessage(String pattern, String channel, String message) { - listener.onMessage(new DefaultMessage(message.getBytes(), channel.getBytes()), pattern.getBytes()); + listener.onMessage(new DefaultMessage(channel.getBytes(), message.getBytes()), pattern.getBytes()); } @Override diff --git a/spring-data-redis/src/test/java/org/springframework/data/keyvalue/redis/listener/PubSubTests.java b/spring-data-redis/src/test/java/org/springframework/data/keyvalue/redis/listener/PubSubTests.java index 1cfd000ec..eec7e5f94 100644 --- a/spring-data-redis/src/test/java/org/springframework/data/keyvalue/redis/listener/PubSubTests.java +++ b/spring-data-redis/src/test/java/org/springframework/data/keyvalue/redis/listener/PubSubTests.java @@ -53,7 +53,7 @@ public class PubSubTests { protected RedisTemplate template; private static Set connFactories = new LinkedHashSet(); - private final BlockingDeque bag = new LinkedBlockingDeque(4); + private final BlockingDeque bag = new LinkedBlockingDeque(99); private final Object handler = new Object() { void handleMessage(String message) { @@ -130,7 +130,25 @@ public class PubSubTests { set.add(bag.poll(1, TimeUnit.SECONDS)); set.add(bag.poll(1, TimeUnit.SECONDS)); + assertTrue(set.contains(payload1)); assertTrue(set.contains(payload2)); } + + @Test + public void testMessageBatch() throws Exception { + + container.addMessageListener(adapter, Arrays.asList(new ChannelTopic(CHANNEL))); + + // wait for the container to start the registration + + int COUNT = 10; + Thread.sleep(500); + for (int i = 0; i < COUNT; i++) { + template.convertAndSend(CHANNEL, "message=" + i); + } + + Thread.sleep(1000); + assertEquals(COUNT, bag.size()); + } } \ No newline at end of file