From da64919b6878ed285dc0609b8b22898472b92a2f Mon Sep 17 00:00:00 2001 From: Costin Leau Date: Wed, 6 Apr 2011 09:16:24 +0300 Subject: [PATCH] Revert "Revert "+ upgrade to Rjc 0.6.4 (snapshot for now)"" This reverts commit 6cf1a58606344d52eedfafe280ab5f9323e98e7e. Upgrade back to RJC 0.6.4 now that it has been released (just in time for M3) --- spring-data-keyvalue-parent/pom.xml | 4 +- spring-data-redis/pom.xml | 4 +- .../rjc/CloseSuppressingRjcConnection.java | 117 ++++++++++++++++++ .../redis/connection/rjc/RjcConnection.java | 14 +-- .../redis/connection/rjc/RjcSubscription.java | 39 +----- 5 files changed, 133 insertions(+), 45 deletions(-) create mode 100644 spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/CloseSuppressingRjcConnection.java diff --git a/spring-data-keyvalue-parent/pom.xml b/spring-data-keyvalue-parent/pom.xml index b0a7cece1..45e6f69bd 100644 --- a/spring-data-keyvalue-parent/pom.xml +++ b/spring-data-keyvalue-parent/pom.xml @@ -314,8 +314,8 @@ org.apache.maven.plugins maven-compiler-plugin - 1.6 - 1.6 + 1.5 + 1.5 -Xlint:all true false diff --git a/spring-data-redis/pom.xml b/spring-data-redis/pom.xml index 2c95966af..df79d0267 100644 --- a/spring-data-redis/pom.xml +++ b/spring-data-redis/pom.xml @@ -16,10 +16,10 @@ "[3.0.0, 4.0.0)" 03122010 1.5.2 - 0.6.3 + 0.6.4 "[1.0.0,2.0.0)" "[1.6, 2.0.0)" - "[0.6.3, 0.6.3]" + "[0.6.4, 0.6.4]" diff --git a/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/CloseSuppressingRjcConnection.java b/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/CloseSuppressingRjcConnection.java new file mode 100644 index 000000000..f5adba611 --- /dev/null +++ b/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/CloseSuppressingRjcConnection.java @@ -0,0 +1,117 @@ +/* + * Copyright 2011 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.data.keyvalue.redis.connection.rjc; + +import java.io.IOException; +import java.net.UnknownHostException; +import java.util.List; + +import org.idevlab.rjc.ds.RedisConnection; +import org.idevlab.rjc.message.RedisNodeSubscriber; +import org.idevlab.rjc.protocol.Protocol.Command; + +/** + * Basic decorator suppressing close() calls to the underlying connection. + * Used for reusing arbitrary connections with {@link RedisNodeSubscriber} without + * resorting to connection pooling. + * + * @author Costin Leau + */ +class CloseSuppressingRjcConnection implements RedisConnection { + + private final RedisConnection delegate; + + /** + * Constructs a new CloseSuppressingRjcConnection instance. + * + * @param delegate + */ + CloseSuppressingRjcConnection(RedisConnection delegate) { + this.delegate = delegate; + } + + public void close() { + // no-op + } + + public void connect() throws UnknownHostException, IOException { + delegate.connect(); + } + + public List getAll() { + return delegate.getAll(); + } + + public String getBulkReply() { + return delegate.getBulkReply(); + } + + public String getHost() { + return delegate.getHost(); + } + + public Long getIntegerReply() { + return delegate.getIntegerReply(); + } + + public List getMultiBulkReply() { + return delegate.getMultiBulkReply(); + } + + public List getObjectMultiBulkReply() { + return delegate.getObjectMultiBulkReply(); + } + + public Object getOne() { + return delegate.getOne(); + } + + public int getPort() { + return delegate.getPort(); + } + + public String getStatusCodeReply() { + return delegate.getStatusCodeReply(); + } + + public int getTimeout() { + return delegate.getTimeout(); + } + + public boolean isConnected() { + return delegate.isConnected(); + } + + public void rollbackTimeout() { + delegate.rollbackTimeout(); + } + + public void sendCommand(Command arg0, byte[]... arg1) { + delegate.sendCommand(arg0, arg1); + } + + public void sendCommand(Command arg0, String... arg1) { + delegate.sendCommand(arg0, arg1); + } + + public void sendCommand(Command arg0) { + delegate.sendCommand(arg0); + } + + public void setTimeoutInfinite() { + delegate.setTimeoutInfinite(); + } +} \ No newline at end of file diff --git a/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcConnection.java b/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcConnection.java index a52af61d5..50d5f78f0 100644 --- a/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcConnection.java +++ b/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcConnection.java @@ -58,7 +58,7 @@ public class RjcConnection implements RedisConnection { SingleDataSource connectionDataSource = new SingleDataSource(connection); session = new SessionFactoryImpl(connectionDataSource).create(); subscriber = new RedisNodeSubscriber(); - subscriber.setDataSource(connectionDataSource); + subscriber.setDataSource(new SingleDataSource(new CloseSuppressingRjcConnection(connection))); client = new Client(connection); this.dbIndex = dbIndex; @@ -79,13 +79,9 @@ public class RjcConnection implements RedisConnection { @Override public void close() throws DataAccessException { isClosed = true; - try { - subscriber.close(); - } catch (Exception ex) { - // ignore - } try { + subscriber.close(); session.close(); } catch (Exception ex) { throw convertRjcAccessException(ex); @@ -2028,8 +2024,9 @@ public class RjcConnection implements RedisConnection { throw new UnsupportedOperationException(); } - subscription = new RjcSubscription(listener, subscriber, client); + subscription = new RjcSubscription(listener, subscriber); subscription.pSubscribe(patterns); + subscriber.runSubscription(); } catch (Exception ex) { throw convertRjcAccessException(ex); @@ -2051,8 +2048,9 @@ public class RjcConnection implements RedisConnection { throw new UnsupportedOperationException(); } - subscription = new RjcSubscription(listener, subscriber, client); + subscription = new RjcSubscription(listener, subscriber); subscription.subscribe(channels); + subscriber.runSubscription(); } catch (Exception ex) { throw convertRjcAccessException(ex); diff --git a/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcSubscription.java b/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcSubscription.java index 64d8bc959..78c5f2277 100644 --- a/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcSubscription.java +++ b/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcSubscription.java @@ -15,7 +15,6 @@ */ package org.springframework.data.keyvalue.redis.connection.rjc; -import org.idevlab.rjc.Client; import org.idevlab.rjc.message.RedisNodeSubscriber; import org.springframework.data.keyvalue.redis.connection.MessageListener; import org.springframework.data.keyvalue.redis.connection.util.AbstractSubscription; @@ -28,62 +27,36 @@ import org.springframework.data.keyvalue.redis.connection.util.AbstractSubscript class RjcSubscription extends AbstractSubscription { private final RedisNodeSubscriber subscriber; - private final Client client; - // rjc does not support subscription while listening - // so we have to handle this ourselves through the client - private volatile boolean subscribed = false; - RjcSubscription(MessageListener listener, RedisNodeSubscriber subscriber, Client client) { + RjcSubscription(MessageListener listener, RedisNodeSubscriber subscriber) { super(listener); this.subscriber = subscriber; subscriber.setMessageListener(new RjcMessageListener(listener)); subscriber.setPMessageListener(new RjcMessageListener(listener)); - this.client = client; } @Override protected void doClose() { - subscribed = false; - client.unsubscribe(); - client.punsubscribe(); - client.rollbackTimeout(); + subscriber.close(); } @Override protected void doPsubscribe(byte[]... patterns) { - String[] pats = RjcUtils.decodeMultiple(patterns); - - if (subscribed) { - client.psubscribe(pats); - } - else { - subscriber.setPatterns(RjcUtils.addArray(subscriber.getPatterns(), pats)); - subscribed = true; - subscriber.subscribe(); - } + subscriber.psubscribe(RjcUtils.decodeMultiple(patterns)); } @Override protected void doPUnsubscribe(boolean all, byte[]... patterns) { - client.punsubscribe(RjcUtils.decodeMultiple(patterns)); + subscriber.punsubscribe(RjcUtils.decodeMultiple(patterns)); } @Override protected void doSubscribe(byte[]... channels) { - String[] chs = RjcUtils.decodeMultiple(channels); - - if (subscribed) { - client.subscribe(chs); - } - else { - subscriber.setPatterns(RjcUtils.addArray(subscriber.getChannels(), chs)); - subscribed = true; - subscriber.subscribe(); - } + subscriber.subscribe(RjcUtils.decodeMultiple(channels)); } @Override protected void doUnsubscribe(boolean all, byte[]... channels) { - client.punsubscribe(RjcUtils.decodeMultiple(channels)); + subscriber.punsubscribe(RjcUtils.decodeMultiple(channels)); } } \ No newline at end of file