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 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 subscriber; + + public OnMessageOnNext(final Subscriber 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 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