diff --git a/pom.xml b/pom.xml
index 51c0883..0ca7c2d 100644
--- a/pom.xml
+++ b/pom.xml
@@ -4,7 +4,7 @@
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
4.0.0
- com.netflix.rxjava
+ io.reactivex
rxjava-redis
0.0.1-SNAPSHOT
@@ -12,16 +12,59 @@
redis.clients
jedis
- 2.1.0
+ 2.7.2
-
-
- com.netflix.rxjava
- rxjava-core
- 0.6.0
+ io.reactivex
+ rxjava
+ 1.0.13
+
+
+ org.slf4j
+ slf4j-api
+ 1.7.7
+
+
+ org.slf4j
+ slf4j-log4j12
+ 1.7.7
+ test
+
+
+ log4j
+ log4j
+ 1.2.17
+ test
+
+
+ org.testng
+ testng
+ 6.8.8
+ test
-
-
+
+
+
+
+ src/test/resources
+
+ **/*.properties
+
+
+
+
+
+ org.apache.maven.plugins
+ maven-compiler-plugin
+ 3.0
+
+
+ 1.7
+ 1.7
+
+
+
+
+
\ No newline at end of file
diff --git a/src/main/java/rx/redis/RedisPoolPubSub.java b/src/main/java/rx/redis/RedisPoolPubSub.java
new file mode 100644
index 0000000..e3e1698
--- /dev/null
+++ b/src/main/java/rx/redis/RedisPoolPubSub.java
@@ -0,0 +1,154 @@
+package rx.redis;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import redis.clients.jedis.Jedis;
+import redis.clients.jedis.JedisPool;
+import redis.clients.jedis.JedisPubSub;
+import redis.clients.jedis.exceptions.JedisConnectionException;
+import rx.Observable;
+import rx.Subscriber;
+import rx.functions.Action0;
+import rx.functions.Func0;
+import rx.subscriptions.Subscriptions;
+
+import java.util.concurrent.Executor;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+
+/**
+ * @author Matteo Moci ( matteo (dot) moci (at) gmail (dot) com )
+ */
+public class RedisPoolPubSub {
+
+ private static final Logger LOGGER = LoggerFactory.getLogger(RedisPoolPubSub.class);
+
+ public static Observable observe(final JedisPool jedisPool, final String channel) {
+
+ return observe(jedisPool, channel, Executors.newSingleThreadExecutor());
+ }
+
+ public static Observable observe(final JedisPool jedisPool, final String channel,
+ final ExecutorService executor) {
+
+ return Observable.defer(new Func0>() {
+ @Override
+ public Observable call() {
+
+ return Observable.create(new RedisObservable(jedisPool, channel, executor));
+ }
+ });
+ }
+
+ private static class RedisObservable implements Observable.OnSubscribe {
+
+ private final JedisPool jedisPool;
+
+ private final String channel;
+
+ private final Executor executor;
+
+ public RedisObservable(final JedisPool jedisPool, final String channel,
+ final Executor executor) {
+
+ this.jedisPool = jedisPool;
+
+ this.channel = channel;
+
+ this.executor = executor;
+
+ }
+
+ @Override
+ public void call(final Subscriber super String> subscriber) {
+
+ final JedisPubSub jedisPubSub = new OnMessageOnNext(subscriber);
+
+ //the resource is returned to the pool by UnsubscribeAction
+ final Jedis jedis;
+ try {
+ jedis = jedisPool.getResource();
+
+ //subscribe
+ executor.execute(new Runnable() {
+ @Override
+ public void run() {
+
+ jedis.subscribe(jedisPubSub, channel);
+
+ //blocked until it's called jedisPubSub.unsubscribe()
+ //it closes jedis in the end
+
+ jedis.close();
+
+ }
+ });
+
+ //unsubscribe jedisPubSub and close hedis when observer is unsubscribed
+ // from http://stackoverflow.com/questions/26695125/how-to-get-notified-of-a-observers-unsubscribe-action-in-a-custom-observable-in
+ subscriber.add(Subscriptions.create(new Action0() {
+ @Override
+ public void call() {
+
+ jedisPubSub.unsubscribe();
+
+ }
+ }));
+
+ } catch (final JedisConnectionException e) {
+ subscriber.onError(e);
+ }
+ }
+
+ private static class OnMessageOnNext extends JedisPubSub {
+
+ private Subscriber super String> subscriber;
+
+ public OnMessageOnNext(final Subscriber super String> subscriber) {
+
+ this.subscriber = subscriber;
+ }
+
+ @Override
+ public void onMessage(final String channel, final String message) {
+
+ if (!subscriber.isUnsubscribed()) {
+ subscriber.onNext(message);
+ }
+
+ }
+
+ @Override
+ public void onPMessage(String pattern, String channel, String message) {
+
+ }
+
+ public void onSubscribe(String channel, int subscribedChannels) {
+ //TODO?
+ // subscriber.onStart();
+
+ }
+
+ @Override
+ public void onUnsubscribe(String channel, int subscribedChannels) {
+
+ //TODO?
+ // subscriber.onCompleted();
+
+ }
+
+ @Override
+ public void onPUnsubscribe(String pattern, int subscribedChannels) {
+
+ }
+
+ @Override
+ public void onPSubscribe(String pattern, int subscribedChannels) {
+
+ }
+
+ }
+ }
+
+}
+
diff --git a/src/main/java/rx/redis/RedisPubSub.java b/src/main/java/rx/redis/RedisPubSub.java
index 7fc131f..dc585a8 100644
--- a/src/main/java/rx/redis/RedisPubSub.java
+++ b/src/main/java/rx/redis/RedisPubSub.java
@@ -1,12 +1,14 @@
package rx.redis;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
import redis.clients.jedis.Jedis;
import redis.clients.jedis.JedisPubSub;
import rx.Observable;
-import rx.Observer;
-import rx.Subscription;
-import rx.util.functions.Func0;
-import rx.util.functions.Func1;
+import rx.Subscriber;
+import rx.functions.Action0;
+import rx.functions.Func0;
+import rx.subscriptions.Subscriptions;
import java.util.concurrent.Executor;
import java.util.concurrent.ExecutorService;
@@ -14,55 +16,75 @@
public final class RedisPubSub {
+ private static final Logger LOGGER = LoggerFactory.getLogger(RedisPubSub.class);
+
public static Observable observe(final Jedis jedis, final String channel) {
+
return observe(jedis, channel, Executors.newSingleThreadExecutor());
}
- public static Observable observe(final Jedis jedis, final String channel, final ExecutorService executor) {
+ public static Observable observe(final Jedis jedis, final String channel,
+ final ExecutorService executor) {
+
return Observable.defer(new Func0>() {
@Override
public Observable call() {
+
return Observable.create(new RedisObservable(jedis, channel, executor));
}
});
}
- private static class RedisObservable implements Func1, Subscription> {
+ private static class RedisObservable implements Observable.OnSubscribe {
private final Jedis jedis;
+
private final String channel;
+
private final Executor executor;
- public RedisObservable(Jedis jedis, String channel, Executor executor) {
+ public RedisObservable(final Jedis jedis, final String channel, final Executor executor) {
+
this.jedis = jedis;
+
this.channel = channel;
+
this.executor = executor;
}
@Override
- public Subscription call(final Observer observer) {
+ public void call(final Subscriber super String> subscriber) {
+
final JedisPubSub pubSub = new JedisPubSub() {
@Override
public void onMessage(String channel, String message) {
- observer.onNext(message);
+
+ subscriber.onNext(message);
}
@Override
- public void onPMessage(String pattern, String channel,
- String message) {
+ public void onPMessage(String pattern, String channel, String message) {
}
public void onSubscribe(String channel, int subscribedChannels) {
+
+ //TODO
+ // subscriber.onStart();
+
}
@Override
public void onUnsubscribe(String channel, int subscribedChannels) {
- }
+ //TODO
+ // subscriber.onCompleted();
+
+ }
@Override
public void onPUnsubscribe(String pattern, int subscribedChannels) {
+
}
@Override
@@ -74,16 +96,22 @@ public void onPSubscribe(String pattern, int subscribedChannels) {
executor.execute(new Runnable() {
@Override
public void run() {
+
jedis.subscribe(pubSub, channel);
+
}
});
- return new Subscription() {
+ // from http://stackoverflow.com/questions/26695125/how-to-get-notified-of-a-observers-unsubscribe-action-in-a-custom-observable-in
+ subscriber.add(Subscriptions.create(new Action0() {
@Override
- public void unsubscribe() {
- pubSub.unsubscribe();
+ public void call() {
+
+ pubSub.unsubscribe(channel);
+
}
- };
+ }));
+
}
}
diff --git a/src/test/java/rx/redis/JedisPoolTestCase.java b/src/test/java/rx/redis/JedisPoolTestCase.java
new file mode 100644
index 0000000..b5534bf
--- /dev/null
+++ b/src/test/java/rx/redis/JedisPoolTestCase.java
@@ -0,0 +1,115 @@
+package rx.redis;
+
+import org.apache.commons.pool2.impl.GenericObjectPoolConfig;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.testng.annotations.Test;
+import redis.clients.jedis.Jedis;
+import redis.clients.jedis.JedisPool;
+import redis.clients.jedis.JedisPubSub;
+import redis.clients.jedis.exceptions.JedisConnectionException;
+
+import static org.testng.Assert.assertNotNull;
+import static org.testng.Assert.assertNull;
+
+/**
+ * @author Matteo Moci ( matteo (dot) moci (at) gmail (dot) com )
+ */
+public class JedisPoolTestCase {
+
+ private static final Logger LOGGER = LoggerFactory.getLogger(JedisPoolTestCase.class);
+
+
+ /* TODO test this script:
+ *
+ * acquire jedis from pool
+ * on exception, throw
+ * register a pubsub on a channel using thread x
+ * on websocketclose, unsubscribe pubsub
+ * the method subscribe terminates, and we can release the jedis.
+ *
+ * */
+
+ @Test (enabled = true)
+ public void testName() throws Exception {
+
+ final GenericObjectPoolConfig poolConfig = new GenericObjectPoolConfig();
+ poolConfig.setMaxIdle(1);
+ poolConfig.setMaxTotal(1);
+ poolConfig.setMinIdle(1);
+ poolConfig.setBlockWhenExhausted(false);
+ poolConfig.setTestOnBorrow(true);
+ poolConfig.setTestWhileIdle(true);
+
+ final JedisPool jedisPool = new JedisPool(poolConfig, "localhost", 6380);
+
+ LOGGER.info("pool before getResource " + jedisPool.getNumActive());
+
+ final Jedis jedis = jedisPool.getResource();
+ LOGGER.info("got resource");
+ LOGGER.info("pool after getResource " + jedisPool.getNumActive());
+
+ final LocalPubSub jedisPubSub = new LocalPubSub();
+
+ new Thread(new Runnable() {
+ @Override
+ public void run() {
+
+ LOGGER.info("subscribing");
+ jedis.subscribe(jedisPubSub, "a-channel");
+ LOGGER.info("subscribe finished, always after unsubscribing");
+
+
+ //from here
+ LOGGER.info("closing, releasing jedis");
+ jedis.close();
+ LOGGER.info("released jedis");
+ //to here
+
+ }
+ }).start();
+
+ LOGGER.info("first thread started");
+
+ Thread.sleep(100L);
+
+ new Thread(new Runnable() {
+ @Override
+ public void run() {
+
+ LOGGER.info("unsubscribing");
+ jedisPubSub.unsubscribe();
+ LOGGER.info("unsubscribed");
+
+
+ }
+ }).start();
+
+ LOGGER.info("second thread started");
+
+ Thread.sleep(5000L);
+
+ LOGGER.info("pool at end " + jedisPool.getNumActive());
+
+ Jedis resource = null;
+ try {
+ resource = jedisPool.getResource();
+ } catch (JedisConnectionException e) {
+ assertNull(e);
+ }
+ assertNotNull(resource);
+
+ }
+
+ private static class LocalPubSub extends JedisPubSub {
+
+ @Override
+ public void onSubscribe(String channel, int subscribedChannels) {
+
+ LOGGER.info(
+ "channel: '" + channel + "' subscribed channels '" + subscribedChannels + "'");
+
+ }
+
+ }
+}
diff --git a/src/test/java/rx/redis/Publisher.java b/src/test/java/rx/redis/Publisher.java
deleted file mode 100644
index 984fd67..0000000
--- a/src/test/java/rx/redis/Publisher.java
+++ /dev/null
@@ -1,22 +0,0 @@
-package rx.redis;
-
-import redis.clients.jedis.Jedis;
-
-import java.util.Arrays;
-
-public class Publisher {
-
- public static void main(String[] args) throws Exception {
- Jedis j = new Jedis("localhost");
- j.connect();
-
- System.out.println("Publishing messages");
- for (String msg : Arrays.asList("one", "two", "three", "four", "five", "six", "seven", "eight", "nine", "ten")) {
- Thread.sleep(1000);
- j.publish("channel", msg);
- }
-
- j.disconnect();
-
- }
-}
diff --git a/src/test/java/rx/redis/RedisPoolPubSubTestCase.java b/src/test/java/rx/redis/RedisPoolPubSubTestCase.java
new file mode 100644
index 0000000..423450f
--- /dev/null
+++ b/src/test/java/rx/redis/RedisPoolPubSubTestCase.java
@@ -0,0 +1,131 @@
+package rx.redis;
+
+import org.apache.commons.pool2.impl.GenericObjectPoolConfig;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.testng.annotations.Test;
+import redis.clients.jedis.Jedis;
+import redis.clients.jedis.JedisPool;
+import rx.Observable;
+import rx.Subscription;
+import rx.functions.Action0;
+import rx.functions.Action1;
+
+import java.util.concurrent.atomic.AtomicBoolean;
+
+import static org.testng.Assert.assertEquals;
+import static org.testng.Assert.assertNotNull;
+import static org.testng.Assert.assertTrue;
+import static org.testng.Assert.fail;
+
+/**
+ * @author Matteo Moci ( matteo (dot) moci (at) gmail (dot) com )
+ */
+public class RedisPoolPubSubTestCase {
+
+ private static final Logger LOGGER = LoggerFactory.getLogger(RedisPubSubTestCase.class);
+
+ private final AtomicBoolean shouldRun;
+
+ public RedisPoolPubSubTestCase() {
+
+ shouldRun = new AtomicBoolean(true);
+ }
+
+ @Test
+ public void testName() throws Exception {
+
+ final Jedis publisherJedis = new Jedis("localhost", 6380);
+
+ final Thread publisher = new Thread(new Runnable() {
+
+ @Override
+ public void run() {
+
+ while (shouldRun.get()) {
+
+ publisherJedis.publish("a-channel", "msg at:'" + System.nanoTime() + "'");
+
+ try {
+ Thread.sleep(1L);
+ } catch (InterruptedException e) {
+ LOGGER.warn("", e);
+ }
+ }
+ }
+ });
+
+ publisher.start();
+
+ final GenericObjectPoolConfig poolConfig = new GenericObjectPoolConfig();
+ poolConfig.setMaxIdle(1);
+ poolConfig.setMaxTotal(1);
+ poolConfig.setBlockWhenExhausted(false);
+
+ final JedisPool subscriberJedisPool = new JedisPool(poolConfig, "localhost", 6380);
+
+ final Observable redisObservable = RedisPoolPubSub.observe(subscriberJedisPool,
+ "a-channel");
+
+ final AtomicBoolean atLeastAMessageWasConsumed = new AtomicBoolean(false);
+
+ final Subscription firstSubscription = redisObservable.subscribe(new Action1() {
+ @Override
+ public void call(final String s) {
+
+ if (!atLeastAMessageWasConsumed.get()) {
+ atLeastAMessageWasConsumed.set(true);
+ }
+ }
+ });
+
+ Thread.sleep(10L);
+
+ final AtomicBoolean secondObservableCalledOnError = new AtomicBoolean(false);
+
+ final Observable secondRedisObservable = RedisPoolPubSub.observe(
+ subscriberJedisPool, "a-channel");
+
+ final Subscription failingSubscription = secondRedisObservable.subscribe(
+ new Action1() {
+ @Override
+ public void call(final String s) {
+ //should never happen
+ fail("should never happen");
+ }
+ }, new Action1() {
+ @Override
+ public void call(final Throwable throwable) {
+
+ //should be called:
+ secondObservableCalledOnError.set(true);
+
+ }
+ }, new Action0() {
+ @Override
+ public void call() {
+ //should never happen
+ fail("should never happen");
+ }
+ });
+
+ shouldRun.set(false);
+
+ firstSubscription.unsubscribe();
+ //need to sleep a bit to work. TODO inspect better
+ Thread.sleep(1L);
+
+ try {
+ final Jedis jedis = subscriberJedisPool.getResource();
+ assertNotNull(jedis);
+ } catch (final Exception e) {
+ fail(e.getMessage());
+ }
+
+ assertTrue(secondObservableCalledOnError.get());
+
+ assertTrue(atLeastAMessageWasConsumed.get());
+
+ assertEquals(subscriberJedisPool.getNumActive(), 1);
+ }
+}
diff --git a/src/test/java/rx/redis/RedisPubSubTestCase.java b/src/test/java/rx/redis/RedisPubSubTestCase.java
new file mode 100644
index 0000000..90f1698
--- /dev/null
+++ b/src/test/java/rx/redis/RedisPubSubTestCase.java
@@ -0,0 +1,69 @@
+package rx.redis;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.testng.annotations.Test;
+import redis.clients.jedis.Jedis;
+import rx.Observable;
+import rx.Subscription;
+
+import java.util.concurrent.atomic.AtomicBoolean;
+
+/**
+ * @author Matteo Moci ( matteo (dot) moci (at) gmail (dot) com )
+ */
+public class RedisPubSubTestCase {
+
+ private static final Logger LOGGER = LoggerFactory.getLogger(RedisPubSubTestCase.class);
+
+ private final AtomicBoolean shouldRun;
+
+ public RedisPubSubTestCase() {
+
+ shouldRun = new AtomicBoolean(true);
+ }
+
+ @Test
+ public void testName() throws Exception {
+
+ final Jedis publisherJedis = new Jedis("localhost", 6380);
+
+ final Thread publisher = new Thread(new Runnable() {
+
+ @Override
+ public void run() {
+
+ while (shouldRun.get()) {
+
+ publisherJedis.publish("a-channel", "msg at:'" + System.nanoTime() + "'");
+
+ try {
+ Thread.sleep(100L);
+ } catch (InterruptedException e) {
+ LOGGER.warn("", e);
+ }
+
+ }
+
+ }
+ });
+
+ publisher.start();
+
+ final Jedis subscriberJedis = new Jedis("localhost", 6380);
+
+ final Observable redisObservable = RedisPubSub.observe(subscriberJedis,
+ "a-channel");
+
+ final Subscription subscribe = redisObservable.subscribe();
+
+ Thread.sleep(1000L);
+
+ shouldRun.set(false);
+
+ subscribe.unsubscribe();
+
+ LOGGER.info("end");
+
+ }
+}
diff --git a/src/test/java/rx/redis/Subscriber.java b/src/test/java/rx/redis/Subscriber.java
deleted file mode 100644
index 06e7cd5..0000000
--- a/src/test/java/rx/redis/Subscriber.java
+++ /dev/null
@@ -1,30 +0,0 @@
-package rx.redis;
-
-import redis.clients.jedis.Jedis;
-import rx.util.functions.Action1;
-import rx.util.functions.Func1;
-
-public class Subscriber {
-
- public static void main(String[] args) throws InterruptedException {
- Jedis j = new Jedis("localhost");
- j.connect();
-
- System.out.println("Subscribing...");
- RedisPubSub.observe(j, "channel")
- .map(new Func1() {
- @Override
- public Integer call(String s) {
- return s.length();
- }
- })
- .subscribe(new Action1() {
- @Override
- public void call(Integer len) {
- System.out.println(len);
- }
- });
-
- }
-
-}
diff --git a/src/test/resources/log4j.properties b/src/test/resources/log4j.properties
new file mode 100644
index 0000000..00c4261
--- /dev/null
+++ b/src/test/resources/log4j.properties
@@ -0,0 +1,5 @@
+datestamp=yyyy-MM-dd HH:mm:ss,SSS/zzz
+log4j.rootLogger=DEBUG, A1
+log4j.appender.A1=org.apache.log4j.ConsoleAppender
+log4j.appender.A1.layout=org.apache.log4j.PatternLayout
+log4j.appender.A1.layout.ConversionPattern=%d{${datestamp}} [%t] %-5p %c %x - %m%n